ConsumerGroupDescribe icon
OutputValues icon
Loop icon
TopicDescribe icon
Log icon
If icon
SlackIncomingWebhook icon
TopicCreatePartitions icon
LoopUntil icon
Fail icon
Schedule icon

Kafka Idle Consumer Partition Guard

Detect Kafka consumer groups with idle members because the topic has too few partitions, add partitions, wait for the rebalance and verify.

Categories
DataInfrastructure

Scaling a Kafka consumer from 2 to 5 replicas does nothing if the topic has 2 partitions. A partition is consumed by one member of a group at a time, so 3 consumers get no assignment and sit idle, while lag keeps growing on the other 2. The autoscaler, the dashboard and the pod count all say the service scaled.

This blueprint compares the active members of each consumer group with the partition count of the topics it reads. It alerts Slack with the number of idle members, and can raise the partition count to the number of members with TopicCreatePartitions. It then waits for the consumers to rebalance, and verifies that every member now has a partition.

This blueprint was created by zkasuran.

Decisions per topic

Action Rule
ADD_PARTITIONS More members than partitions: grow to the member count, capped at max_partitions
PROTECTED The topic is in protected_topics: report only
COMPACTED cleanup.policy includes compact: report only, since keys must stay on their partition
AT_MAX Already at max_partitions: report only

Adding partitions is permanent, and it changes which partition a key maps to. Records with the same key written before and after the change can land on different partitions, so per-key ordering breaks across the change. List keyed topics where ordering matters in protected_topics.

How it works

  1. describe_groups (io.kestra.plugin.kafka.ConsumerGroupDescribe, read-only) returns the members, their assigned partitions and the topics each group reads.
  2. describe_topics (io.kestra.plugin.core.flow.Loop over io.kestra.plugin.kafka.TopicDescribe) reads the partition count and cleanup policy of each topic.
  3. findings compares members with partitions per group and topic, counts the idle members and decides the action. summary logs it.
  4. route (io.kestra.plugin.core.flow.If):
    • balanced: logs it.
    • findings: alert posts to Slack.
      • remediate_gate runs only when auto_remediate is true, dry_run is false, and a topic is ADD_PARTITIONS.
      • remediate (io.kestra.plugin.kafka.TopicCreatePartitions, one topic at a time) sets the new total.
      • verify_topics describes the topics again.
      • wait_rebalance (io.kestra.plugin.core.flow.LoopUntil) polls the groups every 20 seconds until no member is idle.
      • verify_result fails if a topic is not at its target count.
  5. errors alerts Slack when the flow fails, for example on an unreachable broker.

What you get

  • The gap between how far a consumer scaled and how far it can actually parallelize, per group and topic.
  • Partition growth only where it is safe, with keyed, compacted and capped topics reported instead.
  • Proof after the change: the new partition count, and every member busy after the rebalance.

Who it's for

  • Teams that autoscale Kafka consumers on Kubernetes (KEDA, HPA) and see lag that scaling does not fix.
  • Platform teams running shared Kafka clusters with topics created at a few partitions.
  • Teams evaluating the Kafka plugin who want an example of ConsumerGroupDescribe driving TopicCreatePartitions.

Why orchestrate this with Kestra

The mismatch appears when a deployment scales, not when the topic is created. Kestra checks the groups on a schedule and sends the idle consumers to Slack. It grows the topic only behind auto_remediate and dry_run, and only where it is safe, waits for the rebalance with LoopUntil, verifies it, and keeps a history of every partition change.

Tested end to end

On Kestra 2.0.5 OSS against Apache Kafka 3.9.1 (KRaft), with a local HTTP sink standing in for Slack. Consumers were started with kafka-console-consumer.sh --group:

  • 5 consumers in group payments-processor on the 2-partition topic payments-events;
  • 1 consumer in orders-sync;
  • 3 consumers in state-loader on the compacted 1-partition topic user-state.
Run Inputs Result
1 defaults payments-processor: 5 members, 2 partitions, 3 idle, grow to 5. orders-sync not flagged. Alert sent
2 payments-events protected, remediate Reported as PROTECTED. kafka-topics.sh still shows 2 partitions
3 dry run Same report, nothing changed
4 remediate 2 to 5 partitions, verified. After the rebalance, kafka-consumer-groups.sh --members shows 1 partition for each of the 5 members
5 defaults Balanced
6 state-loader, remediate Reported as COMPACTED. user-state still has 1 partition
7 unreachable broker errors alert sent

Prerequisites

  • A Kafka cluster reachable from the Kestra worker, and a principal allowed to describe groups and topics and, for remediation, alter topics.
  • Consumers pick up new partitions at their next metadata refresh (metadata.max.age.ms, 5 minutes by default). wait_rebalance waits up to 8 minutes.
  • Local testing: docker run -d --name kafka -p 9092:9092 apache/kafka:3.9.1, then create a 2-partition topic and start 5 kafka-console-consumer.sh with the same --group.

Secrets

  • SLACK_WEBHOOK_URL: used by alert and the errors block. In Kestra OSS, provide it as the base64-encoded environment variable SECRET_SLACK_WEBHOOK_URL.

Inputs

Input Default Purpose
bootstrap_servers localhost:9092 Brokers
groups 2 example groups Consumer groups to check
max_partitions 48 Upper bound when growing a topic
protected_topics [] Topics never repartitioned
auto_remediate false Allow TopicCreatePartitions
dry_run true Must be false for a change to apply

Quick start

  1. Add the SLACK_WEBHOOK_URL secret.
  2. Import the flow and set bootstrap_servers, groups, and protected_topics for keyed topics where per-key ordering matters.
  3. Run it once and read the idle members per group.
  4. Run it with auto_remediate: true and dry_run: false.
  5. Enable the every_15_minutes trigger, which always runs in safe mode.

How to extend

  • Use the maximum replica count of the consumer's autoscaler instead of the current member count as the target.
  • Discover the groups with ConsumerGroupList instead of listing them.
  • Also flag the opposite case: few members and many partitions with growing lag, a sign the consumer should scale out.
  • Ask for approval with a Pause before growing topics read by several groups.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.