diff --git a/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-configmap.yaml b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-configmap.yaml new file mode 100644 index 0000000..303a9db --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-configmap.yaml @@ -0,0 +1,15 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: alertmanager + labels: + app: alertmanager +data: + config.yml: |- + route: + receiver: 'scdf' + + receivers: + - name: 'scdf' + webhook_configs: + - url: http://alertwebhook:8085/alert diff --git a/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-deployment.yaml b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-deployment.yaml new file mode 100644 index 0000000..c2653aa --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-deployment.yaml @@ -0,0 +1,36 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + labels: + app: alertmanager + name: alertmanager +spec: + selector: + matchLabels: + app: alertmanager + template: + metadata: + labels: + app: alertmanager + spec: + containers: + - name: alertmanager + image: prom/alertmanager:latest + args: + - "--config.file=/etc/alertmanager/config.yml" + - "--storage.path=/alertmanager/" + ports: + - name: alertmanager + containerPort: 9093 + volumeMounts: + - name: alertmanager-config-volume + mountPath: /etc/alertmanager/ + - name: alertmanager-storage-volume + mountPath: /alertmanager/ + + volumes: + - name: alertmanager-config-volume + configMap: + name: alertmanager + - name: alertmanager-storage-volume + emptyDir: {} diff --git a/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-service.yaml b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-service.yaml new file mode 100644 index 0000000..88a9fc8 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/alertmanager/prometheus-alertmanager-service.yaml @@ -0,0 +1,12 @@ +apiVersion: v1 +kind: Service +metadata: + name: alertmanager + labels: + app: alertmanager +spec: + selector: + app: alertmanager + ports: + - port: 9093 + targetPort: 9093 diff --git a/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-deployment.yaml b/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-deployment.yaml new file mode 100644 index 0000000..efeab7a --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-deployment.yaml @@ -0,0 +1,44 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: alertwebhook + labels: + app: alertwebhook +spec: + selector: + matchLabels: + app: alertwebhook + replicas: 1 + template: + metadata: + labels: + app: alertwebhook + spec: + containers: + - image: springcloud/scdf-alert-webhook:latest + name: alertwebhook + imagePullPolicy: Always + ports: + - containerPort: 8085 + resources: + limits: + cpu: 0.5 + memory: 2048Mi + requests: + cpu: 0.5 + memory: 1024Mi + env: + - name: SPRING_CLOUD_DATAFLOW_CLIENT_SERVER_URI + value: 'http://${SCDF_SERVER_SERVICE_HOST}:${SCDF_SERVER_SERVICE_PORT}' + - name: SERVER_PORT + value: '8085' + - name: SCDF_ALERT_WEBHOOK_SCALE_IN_FACTOR + value: '1' + - name: SCDF_ALERT_WEBHOOK_SCALE_OUT_FACTOR + value: '4' + - name: SCDF_ALERT_WEBHOOK_SCALE_APPLICATION_NAME + value: 'transform' + - name: SCDF_ALERT_WEBHOOK_SCALE_OUT_ALERT_NAME + value: 'HighThroughputDifference' + - name: SCDF_ALERT_WEBHOOK_SCALE_IN_ALERT_NAME + value: 'ZeroThroughputDifference' diff --git a/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-svc.yaml b/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-svc.yaml new file mode 100644 index 0000000..22ea8e8 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/alertwebhook/alertwebhook-svc.yaml @@ -0,0 +1,14 @@ +apiVersion: v1 +kind: Service +metadata: + name: alertwebhook + labels: + app: alertwebhook +spec: + # If you are running k8s on a local dev box, using minikube, or Kubernetes on docker desktop you can use type NodePort instead + type: LoadBalancer + ports: + - port: 8085 + targetPort: 8085 + selector: + app: alertwebhook diff --git a/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-deployment.yaml b/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-deployment.yaml new file mode 100644 index 0000000..6aceeb5 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-deployment.yaml @@ -0,0 +1,44 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: alertwebhook + labels: + app: alertwebhook +spec: + selector: + matchLabels: + app: alertwebhook + replicas: 1 + template: + metadata: + labels: + app: alertwebhook + spec: + containers: + - image: springcloud/scdf-alert-webhook:latest + name: alertwebhook + imagePullPolicy: Always + ports: + - containerPort: 8085 + resources: + limits: + cpu: 0.5 + memory: 2048Mi + requests: + cpu: 0.5 + memory: 1024Mi + env: + - name: SPRING_CLOUD_DATAFLOW_CLIENT_SERVER_URI + value: 'http://${MY_RELEASE_DATA_FLOW_SERVER_SERVICE_HOST}:${MY_RELEASE_DATA_FLOW_SERVER_SERVICE_PORT}' + - name: SERVER_PORT + value: '8085' + - name: SCDF_ALERT_WEBHOOK_SCALE_IN_FACTOR + value: '1' + - name: SCDF_ALERT_WEBHOOK_SCALE_OUT_FACTOR + value: '4' + - name: SCDF_ALERT_WEBHOOK_SCALE_APPLICATION_NAME + value: 'transform' + - name: SCDF_ALERT_WEBHOOK_SCALE_OUT_ALERT_NAME + value: 'HighThroughputDifference' + - name: SCDF_ALERT_WEBHOOK_SCALE_IN_ALERT_NAME + value: 'ZeroThroughputDifference' diff --git a/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-svc.yaml b/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-svc.yaml new file mode 100644 index 0000000..22ea8e8 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/helm/alertwebhook/alertwebhook-svc.yaml @@ -0,0 +1,14 @@ +apiVersion: v1 +kind: Service +metadata: + name: alertwebhook + labels: + app: alertwebhook +spec: + # If you are running k8s on a local dev box, using minikube, or Kubernetes on docker desktop you can use type NodePort instead + type: LoadBalancer + ports: + - port: 8085 + targetPort: 8085 + selector: + app: alertwebhook diff --git a/dataflow-website/recipes/scaling/kubernetes/helm/prometheus/prometheus-configmap.yaml b/dataflow-website/recipes/scaling/kubernetes/helm/prometheus/prometheus-configmap.yaml new file mode 100644 index 0000000..081cb40 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/helm/prometheus/prometheus-configmap.yaml @@ -0,0 +1,71 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: my-release-prometheus-server + labels: + app: prometheus +data: + alert.rules.yml: |- + groups: + - name: scdfrules + rules: + - alert: HighThroughputDifference + expr: avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="time"}[1m])) by(stream_name) - avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="transform"}[1m])) by(stream_name) > 500 + for: 30s + annotations: + summary: "The throughput difference between time and transform for the {{ $labels.stream_name }} stream exceeded the threshold " + description: "The time app throughput is larger than the transform throughput for {{ $labels.stream_name }}" + - alert: ZeroThroughputDifference + expr: avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="time"}[1m])) by(stream_name) - 4 * (avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="transform"}[1m])) by(stream_name)) <= 1 + for: 3m + annotations: + summary: "The throughput difference between time and transform for the {{ $labels.stream_name }} stream is zero" + description: "The time throughput matches the tranformer throughput for {{ $labels.stream_name }}" + + prometheus.yml: |- + global: + scrape_interval: 10s + scrape_timeout: 9s + evaluation_interval: 10s + rule_files: + - /etc/config/alert.rules.yml + alerting: + alertmanagers: + - scheme: http + static_configs: + - targets: + - "alertmanager:9093" + scrape_configs: + - job_name: 'proxied-applications' + metrics_path: '/metrics/connected' + kubernetes_sd_configs: + - role: pod + namespaces: + names: + - default + relabel_configs: + - source_labels: [__meta_kubernetes_pod_label_app] + action: keep + regex: prometheus-proxy + - source_labels: [__meta_kubernetes_pod_container_port_number] + action: keep + regex: 8080 + - job_name: 'proxies' + metrics_path: '/metrics/proxy' + kubernetes_sd_configs: + - role: pod + namespaces: + names: + - default + relabel_configs: + - source_labels: [__meta_kubernetes_pod_label_app] + action: keep + regex: prometheus-proxy + - source_labels: [__meta_kubernetes_pod_container_port_number] + action: keep + regex: 8080 + - action: labelmap + regex: __meta_kubernetes_pod_label_(.+) + - source_labels: [__meta_kubernetes_pod_name] + action: replace + target_label: kubernetes_pod_name diff --git a/dataflow-website/recipes/scaling/kubernetes/prometheus/prometheus-configmap.yaml b/dataflow-website/recipes/scaling/kubernetes/prometheus/prometheus-configmap.yaml new file mode 100644 index 0000000..083c7d0 --- /dev/null +++ b/dataflow-website/recipes/scaling/kubernetes/prometheus/prometheus-configmap.yaml @@ -0,0 +1,71 @@ +apiVersion: v1 +kind: ConfigMap +metadata: + name: prometheus + labels: + app: prometheus +data: + alert.rules.yml: |- + groups: + - name: scdfrules + rules: + - alert: HighThroughputDifference + expr: avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="time"}[1m])) by(stream_name) - avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="transform"}[1m])) by(stream_name) > 100 + for: 30s + annotations: + summary: "The throughput difference between time and transform for the {{ $labels.stream_name }}" + description: "The load-generator throughput is larger than the transform throughput for {{ $labels.stream_name }}" + - alert: ZeroThroughputDifference + expr: avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="time"}[1m])) by(stream_name) - 4 * (avg(irate(spring_integration_send_seconds_count{name!="errorChannel",name!="nullChannel",result="success",type="channel",application_name="transform"}[1m])) by(stream_name)) <= 1 + for: 3m + annotations: + summary: "The throughput difference between time and transform for the {{ $labels.stream_name }} is zero" + description: "The load-generator throughput matches the tranformer throughput for {{ $labels.stream_name }}" + + prometheus.yml: |- + global: + scrape_interval: 10s + scrape_timeout: 9s + evaluation_interval: 10s + rule_files: + - /etc/prometheus/alert.rules.yml + alerting: + alertmanagers: + - scheme: http + static_configs: + - targets: + - "alertmanager:9093" + scrape_configs: + - job_name: 'proxied-applications' + metrics_path: '/metrics/connected' + kubernetes_sd_configs: + - role: pod + namespaces: + names: + - default + relabel_configs: + - source_labels: [__meta_kubernetes_pod_label_app] + action: keep + regex: prometheus-proxy + - source_labels: [__meta_kubernetes_pod_container_port_number] + action: keep + regex: 8080 + - job_name: 'proxies' + metrics_path: '/metrics/proxy' + kubernetes_sd_configs: + - role: pod + namespaces: + names: + - default + relabel_configs: + - source_labels: [__meta_kubernetes_pod_label_app] + action: keep + regex: prometheus-proxy + - source_labels: [__meta_kubernetes_pod_container_port_number] + action: keep + regex: 8080 + - action: labelmap + regex: __meta_kubernetes_pod_label_(.+) + - source_labels: [__meta_kubernetes_pod_name] + action: replace + target_label: kubernetes_pod_name diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/.gitignore b/dataflow-website/recipes/scaling/scdf-alert-webhook/.gitignore new file mode 100644 index 0000000..a2a3040 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/.gitignore @@ -0,0 +1,31 @@ +HELP.md +target/ +!.mvn/wrapper/maven-wrapper.jar +!**/src/main/** +!**/src/test/** + +### STS ### +.apt_generated +.classpath +.factorypath +.project +.settings +.springBeans +.sts4-cache + +### IntelliJ IDEA ### +.idea +*.iws +*.iml +*.ipr + +### NetBeans ### +/nbproject/private/ +/nbbuild/ +/dist/ +/nbdist/ +/.nb-gradle/ +build/ + +### VS Code ### +.vscode/ diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/MavenWrapperDownloader.java b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/MavenWrapperDownloader.java new file mode 100644 index 0000000..d2ddc9c --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/MavenWrapperDownloader.java @@ -0,0 +1,122 @@ +/* + * Copyright 2012-2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import java.net.*; +import java.io.*; +import java.nio.channels.*; +import java.util.Properties; + +public class MavenWrapperDownloader { + + private static final String WRAPPER_VERSION = "0.5.5"; + /** + * Default URL to download the maven-wrapper.jar from, if no 'downloadUrl' is provided. + */ + private static final String DEFAULT_DOWNLOAD_URL = "https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/" + + WRAPPER_VERSION + "/maven-wrapper-" + WRAPPER_VERSION + ".jar"; + + /** + * Path to the maven-wrapper.properties file, which might contain a downloadUrl property to + * use instead of the default one. + */ + private static final String MAVEN_WRAPPER_PROPERTIES_PATH = + ".mvn/wrapper/maven-wrapper.properties"; + + /** + * Path where the maven-wrapper.jar will be saved to. + */ + private static final String MAVEN_WRAPPER_JAR_PATH = + ".mvn/wrapper/maven-wrapper.jar"; + + /** + * Name of the property which should be used to override the default download url for the wrapper. + */ + private static final String PROPERTY_NAME_WRAPPER_URL = "wrapperUrl"; + + public static void main(String args[]) { + System.out.println("- Downloader started"); + File baseDirectory = new File(args[0]); + System.out.println("- Using base directory: " + baseDirectory.getAbsolutePath()); + + // If the maven-wrapper.properties exists, read it and check if it contains a custom + // wrapperUrl parameter. + File mavenWrapperPropertyFile = new File(baseDirectory, MAVEN_WRAPPER_PROPERTIES_PATH); + String url = DEFAULT_DOWNLOAD_URL; + if (mavenWrapperPropertyFile.exists()) { + FileInputStream mavenWrapperPropertyFileInputStream = null; + try { + mavenWrapperPropertyFileInputStream = new FileInputStream(mavenWrapperPropertyFile); + Properties mavenWrapperProperties = new Properties(); + mavenWrapperProperties.load(mavenWrapperPropertyFileInputStream); + url = mavenWrapperProperties.getProperty(PROPERTY_NAME_WRAPPER_URL, url); + } + catch (IOException e) { + System.out.println("- ERROR loading '" + MAVEN_WRAPPER_PROPERTIES_PATH + "'"); + } + finally { + try { + if (mavenWrapperPropertyFileInputStream != null) { + mavenWrapperPropertyFileInputStream.close(); + } + } + catch (IOException e) { + // Ignore ... + } + } + } + System.out.println("- Downloading from: " + url); + + File outputFile = new File(baseDirectory.getAbsolutePath(), MAVEN_WRAPPER_JAR_PATH); + if (!outputFile.getParentFile().exists()) { + if (!outputFile.getParentFile().mkdirs()) { + System.out.println( + "- ERROR creating output directory '" + outputFile.getParentFile().getAbsolutePath() + "'"); + } + } + System.out.println("- Downloading to: " + outputFile.getAbsolutePath()); + try { + downloadFileFromURL(url, outputFile); + System.out.println("Done"); + System.exit(0); + } + catch (Throwable e) { + System.out.println("- Error downloading"); + e.printStackTrace(); + System.exit(1); + } + } + + private static void downloadFileFromURL(String urlString, File destination) throws Exception { + if (System.getenv("MVNW_USERNAME") != null && System.getenv("MVNW_PASSWORD") != null) { + String username = System.getenv("MVNW_USERNAME"); + char[] password = System.getenv("MVNW_PASSWORD").toCharArray(); + Authenticator.setDefault(new Authenticator() { + @Override + protected PasswordAuthentication getPasswordAuthentication() { + return new PasswordAuthentication(username, password); + } + }); + } + URL website = new URL(urlString); + ReadableByteChannel rbc; + rbc = Channels.newChannel(website.openStream()); + FileOutputStream fos = new FileOutputStream(destination); + fos.getChannel().transferFrom(rbc, 0, Long.MAX_VALUE); + fos.close(); + rbc.close(); + } + +} diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.jar b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.jar new file mode 100644 index 0000000..0d5e649 Binary files /dev/null and b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.jar differ diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.properties b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.properties new file mode 100644 index 0000000..7d59a01 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/.mvn/wrapper/maven-wrapper.properties @@ -0,0 +1,2 @@ +distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.6.2/apache-maven-3.6.2-bin.zip +wrapperUrl=https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.5/maven-wrapper-0.5.5.jar diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/README.md b/dataflow-website/recipes/scaling/scdf-alert-webhook/README.md new file mode 100644 index 0000000..275ddae --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/README.md @@ -0,0 +1,27 @@ + +# Auto-scaling adapter for Streaming Data Pipelines + +The `AlertWebHookApplication` is a [Prometheus Alertmanager Webhook Receiver](https://github.com/prometheus/alertmanager) +that listens for pre-configured Prometheus alerts and leverages the SCDF Scale API to scale out or in a preconfigured +stream application. + +Use the `scdf.alert.webhook.scaleOutAlertName` and `scdf.alert.webhook.scaleInAlertName` to configure the names of the scale-out and scale-in alert names. + +On scale-out alert, the `AlertWebHookApplication` increases the number of application (defined by `scdf.alert.webhook.scaleApplicationName` or `application_name` alert label) instances to the `scdf.alert.webhook.scaleOutFactor` count. + +On scale-in alert, the `AlertWebHookApplication` decreases the number of application (defined by `scdf.alert.webhook.scaleApplicationName` or `application_name` alert label) instances to the `scdf.alert.webhook.scaleInFactor` count. + +The stream name should be provided either as `scdf.alert.webhook.scaleStreamName` property or an Alert label called: `stream_name`. + +The application name should be provided either as `scdf.alert.webhook.scaleApplicationName` property or an Alert label called: `application_name`. + +Internally, `AlertWebHookApplication` uses the Data Flow Scale REST API to scale the app instances in platform agnostic way. + +Use the `spring.cloud.dataflow.client.server-uri` property to configure url of the Data Flow server. + +Following video illustrates the overall flow and the role for the AlertWebHookApplication: + + + diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw b/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw new file mode 100755 index 0000000..822f699 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw @@ -0,0 +1,322 @@ +#!/bin/sh +# ---------------------------------------------------------------------------- +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# https://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. +# ---------------------------------------------------------------------------- + +# ---------------------------------------------------------------------------- +# Maven2 Start Up Batch script +# +# Required ENV vars: +# ------------------ +# JAVA_HOME - location of a JDK home dir +# +# Optional ENV vars +# ----------------- +# M2_HOME - location of maven2's installed home dir +# MAVEN_OPTS - parameters passed to the Java VM when running Maven +# e.g. to debug Maven itself, use +# set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 +# MAVEN_SKIP_RC - flag to disable loading of mavenrc files +# ---------------------------------------------------------------------------- + +if [ -z "$MAVEN_SKIP_RC" ]; then + + if [ -f /etc/mavenrc ]; then + . /etc/mavenrc + fi + + if [ -f "$HOME/.mavenrc" ]; then + . "$HOME/.mavenrc" + fi + +fi + +# OS specific support. $var _must_ be set to either true or false. +cygwin=false +darwin=false +mingw=false +case "$(uname)" in +CYGWIN*) cygwin=true ;; +MINGW*) mingw=true ;; +Darwin*) + darwin=true + # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home + # See https://developer.apple.com/library/mac/qa/qa1170/_index.html + if [ -z "$JAVA_HOME" ]; then + if [ -x "/usr/libexec/java_home" ]; then + export JAVA_HOME="$(/usr/libexec/java_home)" + else + export JAVA_HOME="/Library/Java/Home" + fi + fi + ;; +esac + +if [ -z "$JAVA_HOME" ]; then + if [ -r /etc/gentoo-release ]; then + JAVA_HOME=$(java-config --jre-home) + fi +fi + +if [ -z "$M2_HOME" ]; then + ## resolve links - $0 may be a link to maven's home + PRG="$0" + + # need this for relative symlinks + while [ -h "$PRG" ]; do + ls=$(ls -ld "$PRG") + link=$(expr "$ls" : '.*-> \(.*\)$') + if expr "$link" : '/.*' >/dev/null; then + PRG="$link" + else + PRG="$(dirname "$PRG")/$link" + fi + done + + saveddir=$(pwd) + + M2_HOME=$(dirname "$PRG")/.. + + # make it fully qualified + M2_HOME=$(cd "$M2_HOME" && pwd) + + cd "$saveddir" + # echo Using m2 at $M2_HOME +fi + +# For Cygwin, ensure paths are in UNIX format before anything is touched +if $cygwin; then + [ -n "$M2_HOME" ] && + M2_HOME=$(cygpath --unix "$M2_HOME") + [ -n "$JAVA_HOME" ] && + JAVA_HOME=$(cygpath --unix "$JAVA_HOME") + [ -n "$CLASSPATH" ] && + CLASSPATH=$(cygpath --path --unix "$CLASSPATH") +fi + +# For Mingw, ensure paths are in UNIX format before anything is touched +if $mingw; then + [ -n "$M2_HOME" ] && + M2_HOME="$( ( + cd "$M2_HOME" + pwd + ))" + [ -n "$JAVA_HOME" ] && + JAVA_HOME="$( ( + cd "$JAVA_HOME" + pwd + ))" +fi + +if [ -z "$JAVA_HOME" ]; then + javaExecutable="$(which javac)" + if [ -n "$javaExecutable" ] && ! [ "$(expr \"$javaExecutable\" : '\([^ ]*\)')" = "no" ]; then + # readlink(1) is not available as standard on Solaris 10. + readLink=$(which readlink) + if [ ! $(expr "$readLink" : '\([^ ]*\)') = "no" ]; then + if $darwin; then + javaHome="$(dirname \"$javaExecutable\")" + javaExecutable="$(cd \"$javaHome\" && pwd -P)/javac" + else + javaExecutable="$(readlink -f \"$javaExecutable\")" + fi + javaHome="$(dirname \"$javaExecutable\")" + javaHome=$(expr "$javaHome" : '\(.*\)/bin') + JAVA_HOME="$javaHome" + export JAVA_HOME + fi + fi +fi + +if [ -z "$JAVACMD" ]; then + if [ -n "$JAVA_HOME" ]; then + if [ -x "$JAVA_HOME/jre/sh/java" ]; then + # IBM's JDK on AIX uses strange locations for the executables + JAVACMD="$JAVA_HOME/jre/sh/java" + else + JAVACMD="$JAVA_HOME/bin/java" + fi + else + JAVACMD="$(which java)" + fi +fi + +if [ ! -x "$JAVACMD" ]; then + echo "Error: JAVA_HOME is not defined correctly." >&2 + echo " We cannot execute $JAVACMD" >&2 + exit 1 +fi + +if [ -z "$JAVA_HOME" ]; then + echo "Warning: JAVA_HOME environment variable is not set." +fi + +CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher + +# traverses directory structure from process work directory to filesystem root +# first directory with .mvn subdirectory is considered project base directory +find_maven_basedir() { + + if [ -z "$1" ]; then + echo "Path not specified to find_maven_basedir" + return 1 + fi + + basedir="$1" + wdir="$1" + while [ "$wdir" != '/' ]; do + if [ -d "$wdir"/.mvn ]; then + basedir=$wdir + break + fi + # workaround for JBEAP-8937 (on Solaris 10/Sparc) + if [ -d "${wdir}" ]; then + wdir=$( + cd "$wdir/.." + pwd + ) + fi + # end of workaround + done + echo "${basedir}" +} + +# concatenates all lines of a file +concat_lines() { + if [ -f "$1" ]; then + echo "$(tr -s '\n' ' ' <"$1")" + fi +} + +BASE_DIR=$(find_maven_basedir "$(pwd)") +if [ -z "$BASE_DIR" ]; then + exit 1 +fi + +########################################################################################## +# Extension to allow automatically downloading the maven-wrapper.jar from Maven-central +# This allows using the maven wrapper in projects that prohibit checking in binary data. +########################################################################################## +if [ -r "$BASE_DIR/.mvn/wrapper/maven-wrapper.jar" ]; then + if [ "$MVNW_VERBOSE" = true ]; then + echo "Found .mvn/wrapper/maven-wrapper.jar" + fi +else + if [ "$MVNW_VERBOSE" = true ]; then + echo "Couldn't find .mvn/wrapper/maven-wrapper.jar, downloading it ..." + fi + if [ -n "$MVNW_REPOURL" ]; then + jarUrl="$MVNW_REPOURL/io/takari/maven-wrapper/0.5.5/maven-wrapper-0.5.5.jar" + else + jarUrl="https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.5/maven-wrapper-0.5.5.jar" + fi + while IFS="=" read key value; do + case "$key" in wrapperUrl) + jarUrl="$value" + break + ;; + esac + done <"$BASE_DIR/.mvn/wrapper/maven-wrapper.properties" + if [ "$MVNW_VERBOSE" = true ]; then + echo "Downloading from: $jarUrl" + fi + wrapperJarPath="$BASE_DIR/.mvn/wrapper/maven-wrapper.jar" + if $cygwin; then + wrapperJarPath=$(cygpath --path --windows "$wrapperJarPath") + fi + + if command -v wget >/dev/null; then + if [ "$MVNW_VERBOSE" = true ]; then + echo "Found wget ... using wget" + fi + if [ -z "$MVNW_USERNAME" ] || [ -z "$MVNW_PASSWORD" ]; then + wget "$jarUrl" -O "$wrapperJarPath" + else + wget --http-user=$MVNW_USERNAME --http-password=$MVNW_PASSWORD "$jarUrl" -O "$wrapperJarPath" + fi + elif command -v curl >/dev/null; then + if [ "$MVNW_VERBOSE" = true ]; then + echo "Found curl ... using curl" + fi + if [ -z "$MVNW_USERNAME" ] || [ -z "$MVNW_PASSWORD" ]; then + curl -o "$wrapperJarPath" "$jarUrl" -f + else + curl --user $MVNW_USERNAME:$MVNW_PASSWORD -o "$wrapperJarPath" "$jarUrl" -f + fi + + else + if [ "$MVNW_VERBOSE" = true ]; then + echo "Falling back to using Java to download" + fi + javaClass="$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.java" + # For Cygwin, switch paths to Windows format before running javac + if $cygwin; then + javaClass=$(cygpath --path --windows "$javaClass") + fi + if [ -e "$javaClass" ]; then + if [ ! -e "$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.class" ]; then + if [ "$MVNW_VERBOSE" = true ]; then + echo " - Compiling MavenWrapperDownloader.java ..." + fi + # Compiling the Java class + ("$JAVA_HOME/bin/javac" "$javaClass") + fi + if [ -e "$BASE_DIR/.mvn/wrapper/MavenWrapperDownloader.class" ]; then + # Running the downloader + if [ "$MVNW_VERBOSE" = true ]; then + echo " - Running MavenWrapperDownloader.java ..." + fi + ("$JAVA_HOME/bin/java" -cp .mvn/wrapper MavenWrapperDownloader "$MAVEN_PROJECTBASEDIR") + fi + fi + fi +fi +########################################################################################## +# End of extension +########################################################################################## + +export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"} +if [ "$MVNW_VERBOSE" = true ]; then + echo $MAVEN_PROJECTBASEDIR +fi +MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS" + +# For Cygwin, switch paths to Windows format before running java +if $cygwin; then + [ -n "$M2_HOME" ] && + M2_HOME=$(cygpath --path --windows "$M2_HOME") + [ -n "$JAVA_HOME" ] && + JAVA_HOME=$(cygpath --path --windows "$JAVA_HOME") + [ -n "$CLASSPATH" ] && + CLASSPATH=$(cygpath --path --windows "$CLASSPATH") + [ -n "$MAVEN_PROJECTBASEDIR" ] && + MAVEN_PROJECTBASEDIR=$(cygpath --path --windows "$MAVEN_PROJECTBASEDIR") +fi + +# Provide a "standardized" way to retrieve the CLI args that will +# work with both Windows and non-Windows executions. +MAVEN_CMD_LINE_ARGS="$MAVEN_CONFIG $@" +export MAVEN_CMD_LINE_ARGS + +WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain + +exec "$JAVACMD" \ + $MAVEN_OPTS \ + -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \ + "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \ + ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@" diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw.cmd b/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw.cmd new file mode 100644 index 0000000..84d60ab --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/mvnw.cmd @@ -0,0 +1,182 @@ +@REM ---------------------------------------------------------------------------- +@REM Licensed to the Apache Software Foundation (ASF) under one +@REM or more contributor license agreements. See the NOTICE file +@REM distributed with this work for additional information +@REM regarding copyright ownership. The ASF licenses this file +@REM to you under the Apache License, Version 2.0 (the +@REM "License"); you may not use this file except in compliance +@REM with the License. You may obtain a copy of the License at +@REM +@REM https://www.apache.org/licenses/LICENSE-2.0 +@REM +@REM Unless required by applicable law or agreed to in writing, +@REM software distributed under the License is distributed on an +@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +@REM KIND, either express or implied. See the License for the +@REM specific language governing permissions and limitations +@REM under the License. +@REM ---------------------------------------------------------------------------- + +@REM ---------------------------------------------------------------------------- +@REM Maven2 Start Up Batch script +@REM +@REM Required ENV vars: +@REM JAVA_HOME - location of a JDK home dir +@REM +@REM Optional ENV vars +@REM M2_HOME - location of maven2's installed home dir +@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands +@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a key stroke before ending +@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven +@REM e.g. to debug Maven itself, use +@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 +@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files +@REM ---------------------------------------------------------------------------- + +@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on' +@echo off +@REM set title of command window +title %0 +@REM enable echoing by setting MAVEN_BATCH_ECHO to 'on' +@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO% + +@REM set %HOME% to equivalent of $HOME +if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%") + +@REM Execute a user defined script before this one +if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre +@REM check for pre script, once with legacy .bat ending and once with .cmd ending +if exist "%HOME%\mavenrc_pre.bat" call "%HOME%\mavenrc_pre.bat" +if exist "%HOME%\mavenrc_pre.cmd" call "%HOME%\mavenrc_pre.cmd" +:skipRcPre + +@setlocal + +set ERROR_CODE=0 + +@REM To isolate internal variables from possible post scripts, we use another setlocal +@setlocal + +@REM ==== START VALIDATION ==== +if not "%JAVA_HOME%" == "" goto OkJHome + +echo. +echo Error: JAVA_HOME not found in your environment. >&2 +echo Please set the JAVA_HOME variable in your environment to match the >&2 +echo location of your Java installation. >&2 +echo. +goto error + +:OkJHome +if exist "%JAVA_HOME%\bin\java.exe" goto init + +echo. +echo Error: JAVA_HOME is set to an invalid directory. >&2 +echo JAVA_HOME = "%JAVA_HOME%" >&2 +echo Please set the JAVA_HOME variable in your environment to match the >&2 +echo location of your Java installation. >&2 +echo. +goto error + +@REM ==== END VALIDATION ==== + +:init + +@REM Find the project base dir, i.e. the directory that contains the folder ".mvn". +@REM Fallback to current working directory if not found. + +set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR% +IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir + +set EXEC_DIR=%CD% +set WDIR=%EXEC_DIR% +:findBaseDir +IF EXIST "%WDIR%"\.mvn goto baseDirFound +cd .. +IF "%WDIR%"=="%CD%" goto baseDirNotFound +set WDIR=%CD% +goto findBaseDir + +:baseDirFound +set MAVEN_PROJECTBASEDIR=%WDIR% +cd "%EXEC_DIR%" +goto endDetectBaseDir + +:baseDirNotFound +set MAVEN_PROJECTBASEDIR=%EXEC_DIR% +cd "%EXEC_DIR%" + +:endDetectBaseDir + +IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig + +@setlocal EnableExtensions EnableDelayedExpansion +for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a +@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS% + +:endReadAdditionalConfig + +SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe" +set WRAPPER_JAR="%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.jar" +set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain + +set DOWNLOAD_URL="https://repo.maven.apache.org/maven2/io/takari/maven-wrapper/0.5.5/maven-wrapper-0.5.5.jar" + +FOR /F "tokens=1,2 delims==" %%A IN ("%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.properties") DO ( + IF "%%A"=="wrapperUrl" SET DOWNLOAD_URL=%%B +) + +@REM Extension to allow automatically downloading the maven-wrapper.jar from Maven-central +@REM This allows using the maven wrapper in projects that prohibit checking in binary data. +if exist %WRAPPER_JAR% ( + if "%MVNW_VERBOSE%" == "true" ( + echo Found %WRAPPER_JAR% + ) +) else ( + if not "%MVNW_REPOURL%" == "" ( + SET DOWNLOAD_URL="%MVNW_REPOURL%/io/takari/maven-wrapper/0.5.5/maven-wrapper-0.5.5.jar" + ) + if "%MVNW_VERBOSE%" == "true" ( + echo Couldn't find %WRAPPER_JAR%, downloading it ... + echo Downloading from: %DOWNLOAD_URL% + ) + + powershell -Command "&{"^ + "$webclient = new-object System.Net.WebClient;"^ + "if (-not ([string]::IsNullOrEmpty('%MVNW_USERNAME%') -and [string]::IsNullOrEmpty('%MVNW_PASSWORD%'))) {"^ + "$webclient.Credentials = new-object System.Net.NetworkCredential('%MVNW_USERNAME%', '%MVNW_PASSWORD%');"^ + "}"^ + "[Net.ServicePointManager]::SecurityProtocol = [Net.SecurityProtocolType]::Tls12; $webclient.DownloadFile('%DOWNLOAD_URL%', '%WRAPPER_JAR%')"^ + "}" + if "%MVNW_VERBOSE%" == "true" ( + echo Finished downloading %WRAPPER_JAR% + ) +) +@REM End of extension + +@REM Provide a "standardized" way to retrieve the CLI args that will +@REM work with both Windows and non-Windows executions. +set MAVEN_CMD_LINE_ARGS=%* + +%MAVEN_JAVA_EXE% %JVM_CONFIG_MAVEN_PROPS% %MAVEN_OPTS% %MAVEN_DEBUG_OPTS% -classpath %WRAPPER_JAR% "-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" %WRAPPER_LAUNCHER% %MAVEN_CONFIG% %* +if ERRORLEVEL 1 goto error +goto end + +:error +set ERROR_CODE=1 + +:end +@endlocal & set ERROR_CODE=%ERROR_CODE% + +if not "%MAVEN_SKIP_RC%" == "" goto skipRcPost +@REM check for post script, once with legacy .bat ending and once with .cmd ending +if exist "%HOME%\mavenrc_post.bat" call "%HOME%\mavenrc_post.bat" +if exist "%HOME%\mavenrc_post.cmd" call "%HOME%\mavenrc_post.cmd" +:skipRcPost + +@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on' +if "%MAVEN_BATCH_PAUSE%" == "on" pause + +if "%MAVEN_TERMINATE_CMD%" == "on" exit %ERROR_CODE% + +exit /B %ERROR_CODE% diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/pom.xml b/dataflow-website/recipes/scaling/scdf-alert-webhook/pom.xml new file mode 100644 index 0000000..ecec9a5 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/pom.xml @@ -0,0 +1,115 @@ + + + 4.0.0 + + org.springframework.boot + spring-boot-starter-parent + 2.2.1.RELEASE + + + io.spring.cloud.dataflow.alert.webhook + scdf-alert-webhook + 0.0.1-SNAPSHOT + scdf-alert-webhook + Demo project for Spring Boot + + + 1.8 + + + + + org.springframework.boot + spring-boot-starter-web + + + + org.springframework.cloud + spring-cloud-dataflow-rest-client + 2.3.0.BUILD-SNAPSHOT + + + + org.springframework.boot + spring-boot-starter-test + test + + + org.junit.vintage + junit-vintage-engine + + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + io.fabric8 + docker-maven-plugin + 0.30.0 + + + + springcloud/${project.artifactId} + + springcloud/openjdk + + latest + ${project.version} + + + /tmp + + + + java + -Djava.security.egd=file:/dev/./urandom + -Dlogging.path=/var/log/scdf-alert-webhook + -jar + /maven/scdf-alert-webhook.jar + + + + assembly.xml + + + + + + + + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot + + true + + + + spring-milestones + Spring Milestones + https://repo.spring.io/libs-milestone-local + + false + + + + spring-releases + Spring Releases + https://repo.spring.io/release + + false + + + + diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/docker/assembly.xml b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/docker/assembly.xml new file mode 100644 index 0000000..d0a0769 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/docker/assembly.xml @@ -0,0 +1,15 @@ + + scdf-alert-webhook + + + + io.spring.cloud.dataflow.alert.webhook:scdf-alert-webhook + + . + scdf-alert-webhook.jar + + + diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookApplication.java b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookApplication.java new file mode 100644 index 0000000..2072812 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookApplication.java @@ -0,0 +1,147 @@ +/* + * Copyright 2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.spring.cloud.dataflow.alert.webhook; + +import java.util.Collections; +import java.util.List; +import java.util.Map; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.SpringApplication; +import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.dataflow.rest.client.DataFlowOperations; +import org.springframework.util.CollectionUtils; +import org.springframework.util.StringUtils; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestMethod; +import org.springframework.web.bind.annotation.RestController; + +/** + * Elastic, Stream auto-scaling adapter for Spring Cloud DataFlow. + * + * The {@link AlertWebHookApplication} is configured as a Prometheus Alertmanager Webhook Receiver (https://github.com/prometheus/alertmanager), + * that listens for pre-configured Prometheus alerts and leverages the SCDF Scale API to scale out or in a preconfigured + * stream application. + * + * The {@link AlertWebHookApplication} listens for scale-out and scale-in alerts, configured by the + * {@link AlertWebHookProperties#getScaleOutAlertName} and the {@link AlertWebHookProperties#getScaleInAlertName}. + * + * On scale-out alert it increases the number of application instances to {@link AlertWebHookProperties#getScaleOutFactor} instances. + * + * On scale-in alert it decreases the number of application instances to {@link AlertWebHookProperties#getScaleInFactor} instances. + * + * The stream name should be provided either as a {@link AlertWebHookProperties#getScaleStreamName} property or an + * Alert label called: 'stream_name' + * + * The application name should be provided either as a {@link AlertWebHookProperties#getScaleApplicationName} property or an + * Alert label called: 'application_name' + * + * Internally the {@link AlertWebHookApplication} uses the Data Flow Scale REST API to scale the app instances in platform agnostic way. + * Use the spring.cloud.dataflow.client.server-uri property to configure url of the Data Flow server. + * + * @author Christian Tzolov + */ + +@SpringBootApplication +@RestController +@EnableConfigurationProperties(AlertWebHookProperties.class) +public class AlertWebHookApplication { + + private static final Logger logger = LoggerFactory.getLogger(AlertWebHookApplication.class); + + public static final String FIRING_ALERT = "firing"; + public static final String ALERT_NAME = "alertname"; + public static final String ALERT_LABELS = "labels"; + public static final String ALERT_STATUS = "status"; + public static final String ALERTS = "alerts"; + public static final String STREAM_NAME = "stream_name"; + public static final String APPLICATION_NAME = "application_name"; + + private AlertWebHookProperties properties; + + // Use the spring.cloud.dataflow.client.server-uri property to configure url of the Data Flow server. + private final DataFlowOperations dataFlowOperations; + + @Autowired + public AlertWebHookApplication(DataFlowOperations dataFlowOperations, AlertWebHookProperties properties) { + this.properties = properties; + this.dataFlowOperations = dataFlowOperations; + } + + /** + * This method is called by the AlertManager passing a JSON payload as documented here: + * https://prometheus.io/docs/alerting/configuration/#webhook_config + */ + @RequestMapping(value = "/alert", method = RequestMethod.POST) + public void alertMessage(@RequestBody Map payload) { + + // Extract all alerts received with this message. + List> alerts = (List>) payload.get(ALERTS); + + if (CollectionUtils.isEmpty(alerts)) { + return; // nothing to do + } + + for (Map alert : alerts) { + String alertStatus = (String) alert.get(ALERT_STATUS); + + // Act only upon "firing" alert types + if (alertStatus.equalsIgnoreCase(FIRING_ALERT)) { + + Map labels = (Map) alert.get(ALERT_LABELS); + String alertName = labels.get(ALERT_NAME); // automatically set by the alert manager. + + String applicationName = StringUtils.isEmpty(this.properties.getScaleApplicationName()) ? + labels.get(APPLICATION_NAME) : this.properties.getScaleApplicationName(); + String streamName = StringUtils.isEmpty(this.properties.getScaleStreamName()) ? + labels.get(STREAM_NAME) : this.properties.getScaleStreamName(); + + if (alertName.equalsIgnoreCase(this.properties.getScaleOutAlertName())) { + logger.info(String.format("Scale Out: %s, %s to %s", streamName, applicationName, + this.properties.getScaleOutFactor())); + this.scale(streamName, applicationName, this.properties.getScaleOutFactor()); + } + else if (alertName.equalsIgnoreCase(this.properties.getScaleInAlertName())) { + logger.info(String.format("Scale In: %s, %s to %s", streamName, applicationName, + this.properties.getScaleOutFactor())); + this.scale(streamName, applicationName, this.properties.getScaleInFactor()); + } + else { + logger.warn(String.format("Fired, non-scale alert: %s with labels: %s", alertName, labels)); + } + } + else { + logger.warn(String.format("Non Firing Alert: %s", alerts)); + } + } + } + + // Spring Cloud Data Flow Scale REST API + private void scale(String streamName, String appName, int scale) { + this.dataFlowOperations.streamOperations().scaleApplicationInstances(streamName, appName, scale, + Collections.emptyMap()); + } + + public static void main(String[] args) { + SpringApplication.run(AlertWebHookApplication.class, args); + } +} diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookProperties.java b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookProperties.java new file mode 100644 index 0000000..40071d0 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/java/io/spring/cloud/dataflow/alert/webhook/AlertWebHookProperties.java @@ -0,0 +1,110 @@ +/* + * Copyright 2019 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.spring.cloud.dataflow.alert.webhook; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * @author Christian Tzolov + */ +@ConfigurationProperties("scdf.alert.webhook") +public class AlertWebHookProperties { + + /** + * Number of application instances to scale out to in case of scale out alert event. + */ + private int scaleOutFactor = 3; + + /** + * Number of application instances to scale in to in case of scale in alert event. + */ + private int scaleInFactor = 1; + + /** + * Name of the stream containing the application to be scaled. + * If scaleStreamName is empty then an attempt will be made to resolve the stream name from alert's label + * called 'stream_name'. + */ + private String scaleStreamName; + + /** + * Stream application name (as appears in the stream definition) to scale out or in. + * If scaleApplicationName is empty then an attempt will be made to resolve the stream name from alert's label + * called 'application_name'. + */ + private String scaleApplicationName; + + /** + * Name of the Alert used to trigger that application scale out. + * The alert is defined in the Prometheus alert.rules.yml configuration. + */ + private String scaleOutAlertName; + + /** + * Name of the Alert used to trigger that application scale in. + * The alert is defined in the Prometheus alert.rules.yml configuration. + */ + private String scaleInAlertName; + + public int getScaleOutFactor() { + return scaleOutFactor; + } + + public void setScaleOutFactor(int scaleOutFactor) { + this.scaleOutFactor = scaleOutFactor; + } + + public String getScaleOutAlertName() { + return scaleOutAlertName; + } + + public void setScaleOutAlertName(String scaleOutAlertName) { + this.scaleOutAlertName = scaleOutAlertName; + } + + public String getScaleInAlertName() { + return scaleInAlertName; + } + + public void setScaleInAlertName(String scaleInAlertName) { + this.scaleInAlertName = scaleInAlertName; + } + + public int getScaleInFactor() { + return scaleInFactor; + } + + public void setScaleInFactor(int scaleInFactor) { + this.scaleInFactor = scaleInFactor; + } + + public String getScaleStreamName() { + return scaleStreamName; + } + + public void setScaleStreamName(String scaleStreamName) { + this.scaleStreamName = scaleStreamName; + } + + public String getScaleApplicationName() { + return scaleApplicationName; + } + + public void setScaleApplicationName(String scaleApplicationName) { + this.scaleApplicationName = scaleApplicationName; + } +} diff --git a/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/resources/application.properties b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/resources/application.properties new file mode 100644 index 0000000..e7ceb33 --- /dev/null +++ b/dataflow-website/recipes/scaling/scdf-alert-webhook/src/main/resources/application.properties @@ -0,0 +1,3 @@ +# Default properties +# spring.cloud.dataflow.client.server-uri=http:// +server.port=8085