id: influxdb-metric-anomaly-escalation
namespace: company.team
inputs:
- id: influxdb_url
type: STRING
defaults: "http://influxdb:8086"
description: InfluxDB server base URL.
- id: influxdb_org
type: STRING
defaults: "my-org"
description: InfluxDB organization name.
- id: influxdb_bucket
type: STRING
defaults: "metrics"
description: InfluxDB bucket storing time-series metrics.
- id: seed_demo_data
type: BOOLEAN
defaults: true
description: Whether to seed demo host CPU metrics before running anomaly detection.
- id: lookback
type: SELECT
defaults: "1h"
values:
- "1h"
- "6h"
- "24h"
description: Time range lookback window for querying InfluxDB metrics.
- id: warning_z
type: SELECT
defaults: "2.0"
values:
- "2.0"
- "2.5"
- "3.0"
description: Minimum Z-score threshold to flag a WARNING metric anomaly.
- id: critical_z
type: SELECT
defaults: "3.5"
values:
- "3.5"
- "4.0"
- "5.0"
description: Minimum Z-score threshold to escalate to a CRITICAL metric anomaly.
triggers:
- id: hourly_schedule
type: io.kestra.plugin.core.trigger.Schedule
cron: "0 * * * *"
disabled: true
inputs:
seed_demo_data: false
tasks:
- id: seed_demo_data_if
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.seed_demo_data }}"
then:
- id: seed_metric_data
type: io.kestra.plugin.influxdb.Write
connection:
url: "{{ inputs.influxdb_url }}"
token: "{{ secret('INFLUXDB_TOKEN') }}"
org: "{{ inputs.influxdb_org }}"
bucket: "{{ inputs.influxdb_bucket }}"
source: |
cpu_usage,host=host-a usage_user=48.0 {{ now() | dateAdd(-12, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=48.0 {{ now() | dateAdd(-11, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=49.0 {{ now() | dateAdd(-10, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=49.0 {{ now() | dateAdd(-9, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=50.0 {{ now() | dateAdd(-8, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=50.0 {{ now() | dateAdd(-7, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=50.0 {{ now() | dateAdd(-6, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=51.0 {{ now() | dateAdd(-5, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=51.0 {{ now() | dateAdd(-4, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=52.0 {{ now() | dateAdd(-3, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=52.0 {{ now() | dateAdd(-2, 'MINUTES') | timestampNano }}
cpu_usage,host=host-a usage_user=51.2 {{ now() | dateAdd(-1, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=48.0 {{ now() | dateAdd(-12, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=48.0 {{ now() | dateAdd(-11, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=49.0 {{ now() | dateAdd(-10, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=49.0 {{ now() | dateAdd(-9, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=50.0 {{ now() | dateAdd(-8, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=50.0 {{ now() | dateAdd(-7, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=50.0 {{ now() | dateAdd(-6, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=51.0 {{ now() | dateAdd(-5, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=51.0 {{ now() | dateAdd(-4, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=52.0 {{ now() | dateAdd(-3, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=52.0 {{ now() | dateAdd(-2, 'MINUTES') | timestampNano }}
cpu_usage,host=host-b usage_user=49.5 {{ now() | dateAdd(-1, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=48.0 {{ now() | dateAdd(-12, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=48.0 {{ now() | dateAdd(-11, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=49.0 {{ now() | dateAdd(-10, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=49.0 {{ now() | dateAdd(-9, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=50.0 {{ now() | dateAdd(-8, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=50.0 {{ now() | dateAdd(-7, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=50.0 {{ now() | dateAdd(-6, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=51.0 {{ now() | dateAdd(-5, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=51.0 {{ now() | dateAdd(-4, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=52.0 {{ now() | dateAdd(-3, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=52.0 {{ now() | dateAdd(-2, 'MINUTES') | timestampNano }}
cpu_usage,host=host-c usage_user=54.0 {{ now() | dateAdd(-1, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=48.0 {{ now() | dateAdd(-12, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=48.0 {{ now() | dateAdd(-11, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=49.0 {{ now() | dateAdd(-10, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=49.0 {{ now() | dateAdd(-9, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=50.0 {{ now() | dateAdd(-8, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=50.0 {{ now() | dateAdd(-7, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=50.0 {{ now() | dateAdd(-6, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=51.0 {{ now() | dateAdd(-5, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=51.0 {{ now() | dateAdd(-4, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=52.0 {{ now() | dateAdd(-3, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=52.0 {{ now() | dateAdd(-2, 'MINUTES') | timestampNano }}
cpu_usage,host=host-d usage_user=58.0 {{ now() | dateAdd(-1, 'MINUTES') | timestampNano }}
- id: query_metrics
type: io.kestra.plugin.influxdb.FluxQuery
connection:
url: "{{ inputs.influxdb_url }}"
token: "{{ secret('INFLUXDB_TOKEN') }}"
org: "{{ inputs.influxdb_org }}"
fetchType: STORE
query: |
from(bucket: "{{ inputs.influxdb_bucket }}")
|> range(start: -{{ inputs.lookback }})
|> filter(fn: (r) => r._measurement == "cpu_usage" and r._field == "usage_user")
|> yield()
- id: ion_to_csv
type: io.kestra.plugin.serdes.csv.IonToCsv
from: "{{ outputs.query_metrics.uri }}"
- id: detect_anomalies
type: io.kestra.plugin.jdbc.duckdb.Query
fetchType: FETCH
inputFiles:
metrics.csv: "{{ outputs.ion_to_csv.uri }}"
sql: |
WITH ranked AS (
SELECT
host,
_value,
_time,
ROW_NUMBER() OVER (PARTITION BY host ORDER BY _time DESC) AS rn
FROM read_csv_auto('{{ workingDir }}/metrics.csv', header=True)
),
latest_points AS (
SELECT host, _value AS latest_val, _time AS latest_time
FROM ranked
WHERE rn = 1
),
baseline_stats AS (
SELECT
host,
AVG(_value) AS mean_val,
STDDEV_SAMP(_value) AS stddev_val
FROM ranked
WHERE rn > 1
GROUP BY host
),
scored AS (
SELECT
l.host,
ROUND(l.latest_val, 2) AS latest_val,
ROUND(b.mean_val, 2) AS mean_val,
ROUND(b.stddev_val, 2) AS stddev_val,
ROUND(ABS(l.latest_val - b.mean_val) / NULLIF(b.stddev_val, 0), 2) AS z_score,
l.latest_time
FROM latest_points l
JOIN baseline_stats b ON l.host = b.host
)
SELECT
host,
latest_val,
mean_val,
stddev_val,
z_score,
CASE
WHEN z_score >= CAST('{{ inputs.critical_z }}' AS DOUBLE) THEN 'CRITICAL'
WHEN z_score >= CAST('{{ inputs.warning_z }}' AS DOUBLE) THEN 'WARNING'
ELSE 'NORMAL'
END AS severity,
latest_time
FROM scored
WHERE z_score >= CAST('{{ inputs.warning_z }}' AS DOUBLE)
ORDER BY z_score DESC;
- id: check_anomalies
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.detect_anomalies.size > 0 }}"
then:
- id: process_anomalies
type: io.kestra.plugin.core.flow.Loop
values: "{{ outputs.detect_anomalies.rows }}"
tasks:
- id: route_severity
type: io.kestra.plugin.core.flow.Switch
value: "{{ item.value | jq('.severity') | first }}"
cases:
CRITICAL:
- id: notify_critical
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "CRITICAL Time-Series Anomaly Detected",
"attachments": [
{
"color": "#FF0000",
"fields": [
{"title": "Host", "value": "{{ item.value | jq('.host') | first }}", "short": true},
{"title": "Z-Score", "value": "{{ item.value | jq('.z_score') | first }}", "short": true},
{"title": "Latest Value", "value": "{{ item.value | jq('.latest_val') | first }}%", "short": true},
{"title": "Baseline Mean", "value": "{{ item.value | jq('.mean_val') | first }}%", "short": true},
{"title": "Std Deviation", "value": "{{ item.value | jq('.stddev_val') | first }}", "short": true},
{"title": "Severity", "value": "CRITICAL", "short": true}
]
}
]
}
WARNING:
- id: notify_warning
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "WARNING Time-Series Metric Anomaly",
"attachments": [
{
"color": "#FFCC00",
"fields": [
{"title": "Host", "value": "{{ item.value | jq('.host') | first }}", "short": true},
{"title": "Z-Score", "value": "{{ item.value | jq('.z_score') | first }}", "short": true},
{"title": "Latest Value", "value": "{{ item.value | jq('.latest_val') | first }}%", "short": true},
{"title": "Baseline Mean", "value": "{{ item.value | jq('.mean_val') | first }}%", "short": true},
{"title": "Severity", "value": "WARNING", "short": true}
]
}
]
}
defaults:
- id: log_default_anomaly
type: io.kestra.plugin.core.log.Log
message: "Anomaly detected on host {{ item.value | jq('.host') | first }} with
z-score {{ item.value | jq('.z_score') | first }}."
else:
- id: log_no_anomalies
type: io.kestra.plugin.core.log.Log
message: "No metric anomalies detected above z-score threshold {{
inputs.warning_z }}. All host metrics remain within baseline limits."
errors:
- id: notify_pipeline_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "InfluxDB Metric Anomaly Pipeline Failed",
"attachments": [
{
"color": "#D00000",
"fields": [
{"title": "Execution ID", "value": "{{ execution.id }}", "short": true},
{"title": "Flow ID", "value": "{{ flow.id }}", "short": true},
{"title": "Error Message", "value": "Execution failed during task execution.", "short": false}
]
}
]
}