New to Kestra?
Use blueprints to kickstart your first workflows.
Detect Kafka consumer groups with idle members because the topic has too few partitions, add partitions, wait for the rebalance and verify.
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.
| 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.
describe_groups (io.kestra.plugin.kafka.ConsumerGroupDescribe, read-only) returns the members, their assigned partitions and the topics each group reads.describe_topics (io.kestra.plugin.core.flow.Loop over io.kestra.plugin.kafka.TopicDescribe) reads the partition count and cleanup policy of each topic.findings compares members with partitions per group and topic, counts the idle members and decides the action. summary logs it.route (io.kestra.plugin.core.flow.If):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.errors alerts Slack when the flow fails, for example on an unreachable broker.ConsumerGroupDescribe driving TopicCreatePartitions.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.
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:
payments-processor on the 2-partition topic payments-events;orders-sync;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 |
metadata.max.age.ms, 5 minutes by default). wait_rebalance waits up to 8 minutes.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.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.| 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 |
SLACK_WEBHOOK_URL secret.bootstrap_servers, groups, and protected_topics for keyed topics where per-key ordering matters.auto_remediate: true and dry_run: false.every_15_minutes trigger, which always runs in safe mode.ConsumerGroupList instead of listing them.Pause before growing topics read by several groups.