ConsumerGroupDescribe icon
OutputValues icon
Switch icon
If icon
Set icon
Log icon
Fail icon
Pause icon
ConsumerGroupAlterOffsets icon
SlackIncomingWebhook icon
Schedule icon

Kafka Consumer Group Offset Checkpoint and Rewind

Save Kafka consumer group offsets to KV hourly and rewind a group to a checkpoint with ConsumerGroupAlterOffsets, approval and verification.

Categories
DataInfrastructure

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:

  • picking the wrong group or the wrong time;
  • resetting while the consumers are still running;
  • not checking that the reset actually landed.

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.

How it works

  1. describe (io.kestra.plugin.kafka.ConsumerGroupDescribe, read-only) returns the group state, active members, and committed and end offsets per partition.
  2. 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.
    • A group with no committed offsets fails the run instead of saving an empty checkpoint.
  3. 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.
    • With dry_run: true (the default) nothing else happens.
    • With 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.
  4. errors alerts Slack when a run fails or is refused.

What you get

  • An offset history per consumer group, taken on a schedule and kept in the KV store, so there is always a known-good position to go back to.
  • A rewind plan per partition (current offset, target, messages to replay) before anything changes.
  • Refusals instead of guesses: no checkpoint before the time, a group with live members, or a checkpoint ahead of the group.
  • An approval pause, then a re-describe that proves every partition landed on the checkpoint.

Who it's for

  • Platform and data engineering teams running Kafka consumers whose releases sometimes have to be replayed.
  • On-call engineers who need to undo a bad consumer release without typing --reset-offsets by hand during an incident.
  • Teams evaluating the Kafka plugin who want a worked example of ConsumerGroupDescribe and ConsumerGroupAlterOffsets with gates around them.

Why orchestrate this with Kestra

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.

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. 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

Prerequisites

  • A Kafka cluster reachable from the Kestra worker, and a principal allowed to describe the group and alter its offsets.
  • Stop the group's consumers before a rewind. The flow refuses to run while members are active.
  • Local testing: 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.

Secrets

  • 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.

Inputs

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.

Quick start

  1. Add the SLACK_WEBHOOK_URL secret.
  2. Import the flow and set bootstrap_servers and group_id.
  3. Run it in CHECKPOINT mode, then enable the hourly_checkpoint trigger.
  4. When a release has to be replayed, stop its consumers and run REWIND with rewind_before set to the release time, first with dry_run: true to read the plan.
  5. Run it again with dry_run: false, resume the paused execution to approve, and restart the consumers.

How to extend

  • Checkpoint several groups by making group_id a list and looping over it.
  • Stop and start the consumers from the flow, for example by scaling a Kubernetes deployment to zero before alter and back after verify.
  • Add a rewind_to_offsets mode that takes explicit offsets, for a replay that does not match a checkpoint.
  • Send the plan to the approver in Slack with a link to the paused execution.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.