id: kafka-consumer-lag-alert
namespace: company.team
description: |
Measure a Kafka consumer group's backlog with the broker's own tooling,
and alert Slack the moment lag crosses your threshold, the group stops
committing, or the group was never registered at all.
triggers:
- id: lag_check_schedule
type: io.kestra.plugin.core.trigger.Schedule
description: Check lag every five minutes. Shipped disabled; enable it after
your first successful manual run so the inputs and threshold are proven
before alerts start firing.
cron: "*/5 * * * *"
disabled: true
inputs:
- id: bootstrap_servers
type: STRING
defaults: localhost:9092
description: Bootstrap address of the Kafka cluster, reachable from the task
container (host:port, comma separated for multiple brokers).
- id: consumer_group
type: STRING
defaults: orders-consumer
description: The consumer group whose backlog should be measured.
- id: max_lag
type: INT
defaults: 10000
description: Alert when the sum of lag across the group's partitions reaches or
exceeds this many messages.
- id: kafka_image
type: STRING
defaults: apache/kafka:3.9.0
description: Container image that carries kafka-consumer-groups.sh. Any official
Apache Kafka image works.
tasks:
- id: report_lag
type: io.kestra.plugin.scripts.shell.Commands
description: Ask the broker for the group's committed offsets versus the end of
each partition, and publish the totals — total lag, worst partition,
partition count, active members, and whether the group exists at all —
through Kestra's stdout outputs protocol. The missing-group case is
detected from the tool's message because kafka-consumer-groups.sh exits 0
even when the group is unknown.
containerImage: "{{ inputs.kafka_image }}"
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
networkMode: host
commands:
- |
set -u
OUT=$(/opt/kafka/bin/kafka-consumer-groups.sh --bootstrap-server {{ inputs.bootstrap_servers }} --describe --group {{ inputs.consumer_group }} 2>&1)
printf '%s\n' "$OUT"
if printf '%s' "$OUT" | grep -q "does not exist"; then
printf '::{"outputs":{"group_exists":false,"total_lag":0,"max_partition_lag":0,"partitions":0,"active_members":0}}::\n'
exit 0
fi
printf '%s\n' "$OUT" | awk -v grp="{{ inputs.consumer_group }}" '
$1 == "GROUP" { inside = 1; next }
inside && $1 == grp {
lag = $6 + 0
total += lag
if (lag > max) max = lag
n += 1
if ($7 != "-") members += 1
}
END {
printf "::{\"outputs\":{\"group_exists\":true,\"total_lag\":%.0f,\"max_partition_lag\":%.0f,\"partitions\":%.0f,\"active_members\":%.0f}}::\n", total, max, n, members
}'
- id: check_group_missing
type: io.kestra.plugin.core.flow.If
description: A group that the broker has never heard of means the consumer has
never committed — the check is pointed at the wrong name or the service
never started, and silence would look like a healthy empty backlog.
condition: "{{ outputs.report_lag.vars.group_exists == false }}"
then:
- id: alert_unknown_group
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Name the unknown group and the cluster so the stream owner can fix
the name or start the consumer instead of trusting a quiet dashboard.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Kafka lag check: consumer group {{ inputs.consumer_group }} does not exist on {{ inputs.bootstrap_servers }} — it has never committed an offset. Check the group name or start the consumer. Execution {{ execution.id }}."
}
- id: check_lag_threshold
type: io.kestra.plugin.core.flow.If
description: Sum the lag the report produced and compare it against the alert
threshold — the single number that decides whether the backlog is a
problem yet.
condition: "{{ outputs.report_lag.vars.group_exists == true and
outputs.report_lag.vars.total_lag >= inputs.max_lag }}"
then:
- id: alert_lag_breach
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Post total lag, the worst partition, how many partitions carry it,
and how many consumers are attached, so the responder can tell a slow
consumer from a dead one.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Kafka lag ALERT: group {{ inputs.consumer_group }} on {{ inputs.bootstrap_servers }} is at {{ outputs.report_lag.vars.total_lag }} messages behind (threshold {{ inputs.max_lag }}). Worst partition: {{ outputs.report_lag.vars.max_partition_lag }}. Spread across {{ outputs.report_lag.vars.partitions }} partition(s) with {{ outputs.report_lag.vars.active_members }} active member(s). Execution {{ execution.id }}."
}
- id: check_stalled_consumer
type: io.kestra.plugin.core.flow.If
description: Lag below the threshold is not always fine — if there is any
backlog at all and not a single consumer is attached to the group, nothing
will ever drain it, so escalate on that shape too.
condition: "{{ outputs.report_lag.vars.group_exists == true and
outputs.report_lag.vars.total_lag > 0 and
outputs.report_lag.vars.active_members == 0 }}"
then:
- id: alert_stalled_consumer
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Flag the orphaned backlog — a group with messages waiting and zero
members means the consumer died or scaled to zero, and the lag will
only grow.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Kafka lag ALERT: group {{ inputs.consumer_group }} on {{ inputs.bootstrap_servers }} has {{ outputs.report_lag.vars.total_lag }} message(s) waiting but ZERO active consumers — the backlog will not drain on its own. Restart or scale the consumer. Execution {{ execution.id }}."
}
else:
- id: log_healthy
type: io.kestra.plugin.core.log.Log
description: Record a clean reading so the execution history doubles as a lag
timeline — the trend is as useful as the alerts.
message: "Kafka group {{ inputs.consumer_group }} healthy: {{
outputs.report_lag.vars.total_lag }} lag across {{
outputs.report_lag.vars.partitions }} partition(s), worst {{
outputs.report_lag.vars.max_partition_lag }}, {{
outputs.report_lag.vars.active_members }} active member(s) (threshold
{{ inputs.max_lag }})."
errors:
- id: alert_check_failure
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
description: Alert when the measurement itself fails — an unreachable broker or
a broken tool produces no numbers, and a check that dies quietly is
indistinguishable from a healthy cluster.
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Kafka lag check FAILED in flow {{ flow.id }} (execution {{ execution.id }}) for group {{ inputs.consumer_group }} on {{ inputs.bootstrap_servers }}. Lag was NOT measured — verify the broker address, the Docker network, and that the image in kafka_image is pullable."
}
outputs:
- id: lag_summary
type: JSON
description: 'The measured backlog — e.g. {"group_exists": true, "total_lag": 8,
"max_partition_lag": 8, "partitions": 1, "active_members": 0} — ready for
dashboards or downstream flows.'
value: "{{ outputs.report_lag.vars | toJson }}"