New to Kestra?
Use blueprints to kickstart your first workflows.
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.
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.
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.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.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.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.errors block pages Slack separately if the flow fails outright, so a broken sentinel is never mistaken for a healthy group.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.status output (HEALTHY / WARNING / CRITICAL / UNREACHABLE) other flows can gate on.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.
bootstrap_servers.consumer_group_id with committed offsets on at least one partition.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.
docker run -d --name kafka-broker -p 9092:9092 apache/kafka:3.7.0docker 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.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 1SLACK_WEBHOOK_URL secret, set consumer_group_id to orders-processor, and run the flow once manually to confirm status reads HEALTHY.consumer_lag_schedule once pointed at a real cluster and group.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.max to a group_by(.topic) summary.properties for a secured cluster, sourced from secrets rather than inline.status into a KV store and alert only on state transitions instead of every run.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.status: one of HEALTHY, WARNING, CRITICAL, UNREACHABLE.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.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.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.