New to Kestra?
Use blueprints to kickstart your first workflows.
Save Kafka consumer group offsets to KV hourly and rewind a group to a checkpoint with ConsumerGroupAlterOffsets, approval and verification.
A consumer release goes out with a bug. For two hours it reads messages, commits their offsets and writes garbage downstream. Once the fix is deployed, those messages need processing again. That means moving the group's committed offsets back to where they were before the release. Doing it by hand under pressure with kafka-consumer-groups.sh --reset-offsets is easy to get wrong:
This blueprint takes an hourly checkpoint of a group's committed offsets in the KV store. On demand, it rewinds the group to the newest checkpoint taken before a given time. The rewind shows the per-partition plan, refuses unsafe cases, waits for approval and alters the offsets. It then describes the group again to prove every partition is on the checkpoint.
This blueprint was created by zkasuran.
describe (io.kestra.plugin.kafka.ConsumerGroupDescribe, read-only) returns the group state, active members, and committed and end offsets per partition.CHECKPOINT mode:save_checkpoint (io.kestra.plugin.core.kv.Set) appends the offsets with a timestamp to the KV key kafka_offsets_<group>, and keeps the newest keep_checkpoints entries.REWIND mode:plan picks the newest checkpoint strictly before rewind_before. For each partition it computes the current offset, the target and the messages that will be replayed. log_plan prints the plan.checks refuses when there is no checkpoint before that time, or when the group still has active members (Kafka cannot move a live group's offsets, and the consumers would overwrite them). It also refuses when a partition's checkpoint is ahead of the current offset or unknown to the group.dry_run: true (the default) nothing else happens.dry_run: false:approve (io.kestra.plugin.core.flow.Pause) waits for someone to resume the execution.alter (io.kestra.plugin.kafka.ConsumerGroupAlterOffsets) sets the checkpoint offsets.verify describes the group again, and verify_result fails if any partition differs from the checkpoint.errors alerts Slack when a run fails or is refused.--reset-offsets by hand during an incident.ConsumerGroupDescribe and ConsumerGroupAlterOffsets with gates around them.kafka-consumer-groups.sh --reset-offsets changes offsets, but it does not remember where the group was an hour ago, check that the consumers are stopped, wait for a second person, or prove the result. Kestra adds the schedule that builds the checkpoint history, the KV store that keeps it, the Pause that turns a rewind into an approved change, and an execution history of every rewind with its plan.
On Kestra 2.0.5 OSS against Apache Kafka 3.9.1 (KRaft), with a local HTTP sink standing in for Slack. The demo topic shipments has 3 partitions, and the group shipment-sync had read all 93 messages.
| Run | Inputs | Result |
|---|---|---|
| 1 | CHECKPOINT |
Checkpoint saved: partitions 0, 1 and 2 at offset 31 |
| 2 | (12 new messages consumed) CHECKPOINT |
Second checkpoint: partition 1 at 43 |
| 3 | REWIND, before the second checkpoint, dry run |
Plan: partition 1 from 43 to 31, 12 to replay. Offsets unchanged |
| 4 | REWIND, before any checkpoint |
Refused, Slack alert |
| 5 | REWIND with a console consumer still running |
Refused: 1 active member, Slack alert |
| 6 | REWIND, dry_run: false |
Paused for approval. After resume, offsets altered and verified. kafka-consumer-groups.sh shows partition 1 at 31 with lag 12, and restarting the consumer replays exactly the 12 messages |
docker run -d --name kafka -p 9092:9092 apache/kafka:3.9.1, then consume a topic with kafka-console-consumer.sh --group <name> to create a group.SLACK_WEBHOOK_URL: Slack incoming webhook used by the errors block. In Kestra OSS, provide it as the base64-encoded environment variable SECRET_SLACK_WEBHOOK_URL.| Input | Default | Purpose |
|---|---|---|
mode |
CHECKPOINT |
CHECKPOINT or REWIND |
bootstrap_servers |
localhost:9092 |
Brokers |
group_id |
shipment-sync |
Consumer group |
rewind_before |
none | Rewind target time |
dry_run |
true |
Must be false for offsets to change |
keep_checkpoints |
168 |
One week of hourly checkpoints |
The hourly_checkpoint trigger is shipped disabled. A rewind is always started by hand.
SLACK_WEBHOOK_URL secret.bootstrap_servers and group_id.CHECKPOINT mode, then enable the hourly_checkpoint trigger.REWIND with rewind_before set to the release time, first with dry_run: true to read the plan.dry_run: false, resume the paused execution to approve, and restart the consumers.group_id a list and looping over it.alter and back after verify.rewind_to_offsets mode that takes explicit offsets, for a replay that does not match a checkpoint.