Schedule icon
ConsumerGroupDescribe icon
Return icon
Switch icon
Log icon
SlackIncomingWebhook icon

Kafka Consumer Lag Sentinel

Describe a Kafka consumer group's lag via the AdminClient on a schedule, with no messages consumed, and alert Slack on warning or critical lag.

Categories
DataInfrastructure

A Kafka consumer group can fall behind for minutes before anyone notices, since nothing in Kafka itself pages on lag. This blueprint turns a scheduled, read-only AdminClient call into a Day 2 health gate: io.kestra.plugin.kafka.ConsumerGroupDescribe returns every partition's committed offset, log-end offset, and lag for the group, a single jq query folds that down to the worst lag and classifies it, and the result routes to a quiet log, an informational Slack note, or a Slack page depending on severity.

How it works

  1. describe_consumer_group (io.kestra.plugin.kafka.ConsumerGroupDescribe) calls the Kafka AdminClient's describe-consumer-group and list-offsets APIs for consumer_group_id against bootstrap_servers, bounded by admin_call_timeout. It returns groups, a list where each entry has groupId, state, members, and offsets (one entry per topic/partition with currentOffset, endOffset, and lag). This is a metadata read: no message is consumed, no offset is committed or altered. allowFailure: true keeps a broker or group lookup failure from aborting the flow outright.
  2. classify_consumer_lag (io.kestra.plugin.core.debug.Return) first guards for a missing groups output (UNREACHABLE), then hands the group list to one jq query that folds every partition's lag down to the single worst value with max, and classifies it: no lag value at all becomes UNREACHABLE, at or above critical_lag_threshold becomes CRITICAL, at or above warning_lag_threshold becomes WARNING, otherwise HEALTHY. The thresholds are spliced into the jq expression from the flow's own inputs.
  3. route_by_lag (io.kestra.plugin.core.flow.Switch) branches on that classification: HEALTHY logs a quiet confirmation; WARNING posts an informational Slack note; CRITICAL and UNREACHABLE each page Slack with the thresholds and the execution id; a defaults fallback logs the raw value for any classification outside the four known cases.
  4. log_audit_result (io.kestra.plugin.core.log.Log) writes one audit line per run, independent of which Switch branch fired, re-deriving the worst partition lag for readability so the check's own history stays queryable even if Slack delivery fails.
  5. The flow-level errors block pages Slack separately if the flow fails outright, so a broken sentinel is never mistaken for a healthy group.
  6. The consumer_lag_schedule trigger describes the group every 5 minutes and ships disabled: true, with explicit inputs: overrides for admin_call_timeout and warning_lag_threshold so the scheduled run's behavior is pinned regardless of this flow's own defaults.

What you get

  • A scheduled, read-only lag check for any Kafka consumer group with zero custom AdminClient code.
  • A single status output (HEALTHY / WARNING / CRITICAL / UNREACHABLE) other flows can gate on.
  • Three distinct Slack messages, one per severity, each carrying the thresholds that were crossed.
  • An audit log line on every run, separate from Slack, with the raw worst-partition lag number, so the check's own history survives a webhook outage.

Who it's for

  • Platform and streaming teams running Kafka who have no orchestrator-side lag alert today.
  • On-call engineers who want a backlog-building signal on a consumer group before an SLA on downstream processing is missed.
  • Teams standardizing Day 2 queue monitoring across mixed messaging systems inside Kestra without writing a bespoke lag-exporter dashboard.

Why orchestrate this with Kestra

Tools like Burrow or a Prometheus lag exporter can watch consumer lag, but that is another service to deploy, scrape, and alert from outside your orchestration layer. Kestra keeps the whole check in one place: the schedule, the read-only AdminClient describe call, the jq-based threshold classification, and the Slack call are all one auditable execution history, with the Switch branch ready to grow into real remediation (scaling consumers, opening a ticket) without rewriting the check itself.

Prerequisites

  • A running Kafka cluster reachable at bootstrap_servers.
  • An existing consumer group consumer_group_id with committed offsets on at least one partition.
  • AdminClient read permission on the cluster (describe consumer groups, list offsets); this flow never calls an AdminClient write or delete operation.
  • A Slack incoming webhook for the WARNING, CRITICAL, UNREACHABLE, and flow-failure alerts.

Secrets

  • SLACK_WEBHOOK_URL: Slack incoming webhook used for every alert branch and the flow-failure alert.

In Kestra OSS (no Enterprise secrets backend), secrets are supplied as environment variables prefixed SECRET_, base64-encoded, and read back in flows with {{ secret('NAME') }}, for example SECRET_SLACK_WEBHOOK_URL=$(echo -n 'https://hooks.slack.com/services/<placeholder>' | base64). This keeps the webhook out of the flow YAML but, per Kestra's own documentation, offers no encryption at rest or access control beyond the host environment; use the Enterprise secrets backend for stronger guarantees. If your cluster requires SASL or SSL, add the extra AdminClient properties (for example sasl.jaas.config) under the task's properties map, sourced from additional {{ secret('NAME') }} values rather than inline.

Quick start

  1. Start a local single-broker Kafka cluster in KRaft mode, on a port well away from Kestra's own UI: docker run -d --name kafka-broker -p 9092:9092 apache/kafka:3.7.0
  2. If Kestra is also running in Docker, attach the broker to Kestra's network so the worker can reach it by container name: docker network connect YOUR_KESTRA_NETWORK kafka-broker Then set bootstrap_servers to kafka-broker:9092. From a host-based Kestra install, use localhost:9092 instead.
  3. Create a topic, produce a few records, and commit offsets with a real consumer group so there is something to describe: docker exec kafka-broker /opt/kafka/bin/kafka-topics.sh --create --topic orders --bootstrap-server localhost:9092 docker exec -i kafka-broker /opt/kafka/bin/kafka-console-producer.sh --topic orders --bootstrap-server localhost:9092 <<< "test message" docker exec kafka-broker /opt/kafka/bin/kafka-console-consumer.sh --topic orders --bootstrap-server localhost:9092 --group orders-processor --from-beginning --max-messages 1
  4. Add the SLACK_WEBHOOK_URL secret, set consumer_group_id to orders-processor, and run the flow once manually to confirm status reads HEALTHY.
  5. Enable consumer_lag_schedule once pointed at a real cluster and group.

How to extend

  • Pass groupIds as a list with more than one entry to watch several consumer groups from a single run, then adjust the jq query to classify per group instead of across all of them.
  • Add a per-topic breakdown to the audit log by changing the jq expression from a flat max to a group_by(.topic) summary.
  • Add SASL or SSL AdminClient properties under properties for a secured cluster, sourced from secrets rather than inline.
  • Push status into a KV store and alert only on state transitions instead of every run.

Inputs

  • bootstrap_servers: Kafka broker host:port list the AdminClient connects to.
  • consumer_group_id: the consumer group to describe.
  • admin_call_timeout: ISO-8601 duration cap on the AdminClient describe call.
  • warning_lag_threshold: worst-partition-lag floor at or above which the group is classified WARNING.
  • critical_lag_threshold: worst-partition-lag floor at or above which the group is classified CRITICAL.

Outputs

  • status: one of HEALTHY, WARNING, CRITICAL, UNREACHABLE.

Links

Pitfalls

  • ConsumerGroupDescribe reports lag as of the moment the AdminClient call runs; a group with highly bursty traffic can look healthy between polls while still falling behind briefly in between.
  • A consumer group with zero committed offsets (brand new, or fully caught up with nothing ever consumed) produces an empty offsets list; the jq query treats that the same as a missing group and reports UNREACHABLE, so confirm the group has run at least once before trusting a non-UNREACHABLE result.
  • Setting warning_lag_threshold or critical_lag_threshold too low relative to real throughput produces false alerts; tune both against an observed healthy baseline before enabling the schedule.
  • This task only reads group metadata; it never calls an AdminClient delete, alter-offsets, or reset operation, so there is nothing destructive to default to dry-run for.
See How

New to Kestra?

Use blueprints to kickstart your first workflows.