New to Kestra?
Use blueprints to kickstart your first workflows.
Describe a Kafka topic's partition replicas and in-sync replicas via the AdminClient on a schedule and alert Slack when redundancy drops critically low.
A Kafka partition can silently lose its replication redundancy, one broker restart or disk failure away from becoming unavailable or losing committed data, and nothing in Kafka itself pages on it. This blueprint turns a scheduled, read-only AdminClient call into a Day 2 health gate: io.kestra.plugin.kafka.TopicDescribe returns every partition's leader, assigned replicas, and in-sync replicas (ISR), a single jq query folds that down to the worst ISR count and whether any partition is under-replicated, and the result routes to a quiet log, an informational Slack note, or a Slack page depending on severity.
describe_topic (io.kestra.plugin.kafka.TopicDescribe) calls the Kafka AdminClient's describe-topic API for topic against bootstrap_servers, bounded by admin_call_timeout. It returns topic, partitionCount, replicationFactor, and partitions, a list with one entry per partition containing partition, leader, replicas, and isr. This is a metadata read: no message is consumed, no partition or configuration is altered. allowFailure: true keeps a broker or topic lookup failure from aborting the flow outright.classify_replication_health (io.kestra.plugin.core.debug.Return) first guards for a missing partitions output (UNREACHABLE), then hands the partition list to one jq query that computes the worst (minimum) in-sync-replica count across every partition and whether any partition has fewer in-sync replicas than assigned replicas, and classifies the result: no data at all becomes UNREACHABLE, a worst ISR count at or below critical_min_isr becomes CRITICAL, any under-replicated partition above that floor becomes WARNING, otherwise HEALTHY. The critical floor is spliced into the jq expression from the flow's own input.route_by_replication_health (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 critical floor 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 ISR count alongside the topic's partition count and replication factor 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 fully replicated topic.replication_health_schedule trigger describes the topic every 10 minutes and ships disabled: true, with explicit inputs: overrides for admin_call_timeout and critical_min_isr 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.Under-replicated-partition alerting usually lives in a JMX exporter and a Prometheus rule, another system to deploy and maintain 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 (triggering a partition reassignment job, opening a ticket) without rewriting the check itself.
bootstrap_servers.topic.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:9092SLACK_WEBHOOK_URL secret and run the flow once manually against the single-broker cluster. With replication factor 1 on a single broker, isr equals replicas for every partition, so status should read HEALTHY (a single-broker cluster has no redundancy to lose in the first place, which is itself worth noting in a real deployment).replication_health_schedule once pointed at a real multi-broker cluster and topic.topic input over every topic a team owns using io.kestra.plugin.core.flow.Loop, instead of checking one topic per execution.min/any fold to a list of only the partitions that are under-replicated.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.topic: the topic to describe.admin_call_timeout: ISO-8601 duration cap on the AdminClient describe call.critical_min_isr: in-sync-replica floor at or below which any partition makes the topic CRITICAL.status: one of HEALTHY, WARNING, CRITICAL, UNREACHABLE.isr equal to replicas, since there is nothing else to replicate to; this flow only becomes meaningful on a multi-broker cluster with replicationFactor greater than 1.TopicDescribe reports replication state as of the moment the AdminClient call runs; a brief ISR shrink during a rolling broker restart can resolve before the next scheduled poll and never surface as WARNING or CRITICAL.critical_min_isr equal to or above the topic's actual replicationFactor makes every partition CRITICAL permanently; keep it strictly below the topic's normal replication factor.