Schedule icon
Commands icon
Docker icon
If icon
SlackIncomingWebhook icon
Log icon

Watch a Kafka Consumer Group's Backlog and Alert Before It Becomes an Outage

Measure a Kafka consumer group's lag on a schedule and alert Slack when it crosses your threshold, the group vanishes, or the consumer stops committing.

Categories
DataInfrastructure

Lag is the heartbeat of a streaming pipeline: healthy, it stays flat; dying, it climbs until orders stop processing and someone asks why the dashboard is frozen. This blueprint turns the broker's own answer — kafka-consumer-groups.sh, the tool every Kafka operator already trusts — into a scheduled, branching decision. Each run computes total lag, the worst partition, and how many consumers are actually attached, then picks between three outcomes: over threshold, orphaned backlog, or quiet and healthy.

How it works

  1. report_lag (io.kestra.plugin.scripts.shell.Commands on the io.kestra.plugin.scripts.runner.docker.Docker runner with networkMode: host) runs kafka-consumer-groups.sh --describe inside the official Apache Kafka image, parses the LAG column with awk, and emits group_exists, total_lag, max_partition_lag, partitions, and active_members through the ::{"outputs": ...}:: protocol. A missing group is detected from the tool's own message, because the command exits 0 even when the group is unknown.
  2. check_group_missing (io.kestra.plugin.core.flow.If) catches the group the broker has never heard of — wrong name or consumer never started — and alerts instead of reporting a deceptively clean zero.
  3. check_lag_threshold (io.kestra.plugin.core.flow.If) compares total lag against the max_lag input and posts the backlog, worst partition, and attached-member count to Slack when breached.
  4. check_stalled_consumer (io.kestra.plugin.core.flow.If) escalates a shape thresholds miss: any messages waiting with zero active members means nothing will ever drain them. Under a clean reading, log_healthy writes the numbers to the execution history.
  5. The errors block alerts when the measurement itself fails, and a five-minute Schedule (shipped disabled) drives the cadence once enabled.

What you get

  • A single trustworthy number — total lag — computed the same way every run, straight from the broker.
  • Three distinct alerts instead of one ambiguous page: over threshold, unknown group, orphaned backlog.
  • A lag_summary JSON output (total_lag, max_partition_lag, partitions, active_members) for dashboards or downstream flows.
  • An execution history that doubles as a lag timeline, with healthy runs logged too.

Who it's for

  • SREs and platform teams running order pipelines, CDC streams, or event ingestion where lag turns into customer-visible delay.
  • Streaming teams who want an alert before the backlog is obvious on a dashboard.
  • Anyone who has typed kafka-consumer-groups.sh --describe by hand at 2am and wished it paged them instead.

Why orchestrate this with Kestra

A cron job could echo the lag, but it could not branch three ways on what it found, alert only on the dangerous shapes, self-report its own failures, and keep every reading in a queryable execution history. Kestra supplies the schedule, the conditional logic, the failure channel, and the audit trail — and because the check is a flow, the same alert can later fan out to PagerDuty or open a Jira ticket by adding one task.

Prerequisites

  • A reachable Kafka cluster — the broker's bootstrap address as seen from the Kestra task runner (with networkMode: host, localhost:9092 covers a broker on the same host; containers on a shared Docker network use their service name).
  • Docker available to the Kestra Worker for the task runner.
  • A Slack incoming webhook for alerts.
  • The consumer group you intend to watch must have committed at least one offset (or you will — correctly — get the unknown-group alert).

Secrets

  • SLACK_WEBHOOK_URL: Slack incoming webhook URL for breach, missing-group, stalled-consumer, and failure alerts.

Quick start

  1. Add the Slack webhook secret to your Kestra namespace.
  2. Set bootstrap_servers and consumer_group to your cluster and group.
  3. Run the flow once and read lag_summary (and the logs) to confirm the numbers match kafka-consumer-groups.sh --describe by hand.
  4. Set disabled: false on the lag_check_schedule trigger and tune max_lag to your traffic.

How to extend

  • Add a second threshold input for a "warning" band so Slack gets a heads-up before the page.
  • Loop over several consumer groups by wrapping report_lag and the checks in a ForEach over a groups input.
  • Feed lag_summary into a Kestra dashboard or a nightly trend report by adding an output-consuming subflow.
  • Route the breach branch to PagerDuty or an incident tool by adding a second notification task beside Slack.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.