Schedule icon
TopicDescribe icon
Return icon
Switch icon
Log icon
SlackIncomingWebhook icon

Kafka Topic Replication Sentinel

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.

Categories
DataInfrastructure

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.

How it works

  1. 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.
  2. 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.
  3. 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.
  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 ISR count alongside the topic's partition count and replication factor 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 fully replicated topic.
  6. The 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.

What you get

  • A scheduled, read-only replication health check for any Kafka topic 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 explaining exactly what redundancy loss was detected.
  • An audit log line on every run, separate from Slack, carrying the topic's partition count, replication factor, and worst observed ISR count, 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 replication alert today.
  • On-call engineers who want a redundancy-loss signal on a critical topic before the next broker failure turns it into an outage or data loss event.
  • Teams standardizing Day 2 cluster health monitoring across mixed messaging systems inside Kestra without standing up a separate Cruise Control or JMX-exporter alert for this one signal.

Why orchestrate this with Kestra

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.

Prerequisites

  • A running Kafka cluster reachable at bootstrap_servers.
  • An existing topic named topic.
  • AdminClient read permission on the cluster (describe topics); this flow never calls an AdminClient write, reassignment, 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 to describe: docker exec kafka-broker /opt/kafka/bin/kafka-topics.sh --create --topic orders --bootstrap-server localhost:9092
  4. Add the SLACK_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).
  5. Enable replication_health_schedule once pointed at a real multi-broker cluster and topic.

How to extend

  • Loop this flow's topic input over every topic a team owns using io.kestra.plugin.core.flow.Loop, instead of checking one topic per execution.
  • Add a per-partition breakdown to the audit log by changing the jq expression from a flat min/any fold to a list of only the partitions that are under-replicated.
  • 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.
  • 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.

Outputs

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

Links

Pitfalls

  • A single-broker test cluster (as in Quick start) always reports 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.
  • Setting 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.
  • This task only reads partition metadata; it never calls an AdminClient reassignment, config-alter, or delete operation, so there is nothing destructive to default to dry-run for.
See How

New to Kestra?

Use blueprints to kickstart your first workflows.