Seek & Snapshots
Replay over KubeMQ — Seek to a timestamp or snapshot, replay bounded by MaxSeekReplay, pre-window timestamps clamped to earliest, and snapshot expiry.
Because every topic is backed by a durable, replayable Events Store log gcp.{t}, a
subscription can be rewound. Seek resets a subscription's position to a point in the past —
either a timestamp or a saved snapshot — and replays the topic log from there into the
subscription's queue. This is how you reprocess messages: redeploy a consumer with a bug fix, then
seek the subscription back to before the bad window and let it re-consume.
How Seek works
A Seek against a subscription:
- Resolves the start sequence from the topic log — from a timestamp (the first message at or after that time) or from a snapshot's captured cursor.
- Purges the subscription queue and drops outstanding leases — in-flight
ack_ids become invalid (this is the reset). - Replays the topic log from the start sequence and re-applies the subscription's filter as
it fans the replayed messages back into
gcp.sub.{s}.
Seek(subscription, time | snapshot)
│ resolve start seq from gcp.{t}
▼
purge sub queue + drop leases (in-flight ack_ids now invalid)
│
▼
replay gcp.{t} from start seq ──(re-apply filter)──▶ refill gcp.sub.{s}
│
▼ bounded by MaxSeekReplay (default 1,000,000) → hit cap = WARN, no silent lossA Seek call looks the same in every client; both forms are shown here:
// Seek to a timestamp.
_, _ = subClient.Seek(ctx, &pubsubpb.SeekRequest{
Subscription: subName,
Target: &pubsubpb.SeekRequest_Time{Time: timestamppb.New(cutoff)},
})
// Seek to a snapshot.
_, _ = subClient.Seek(ctx, &pubsubpb.SeekRequest{
Subscription: subName,
Target: &pubsubpb.SeekRequest_Snapshot{Snapshot: snapshotName},
})# Seek to a timestamp.
subscriber.seek(request={"subscription": sub_path, "time": cutoff})
# Seek to a snapshot.
subscriber.seek(request={"subscription": sub_path, "snapshot": snapshot_path})// Seek to a timestamp.
subscriptionAdminClient.seek(SeekRequest.newBuilder()
.setSubscription(subName.toString()).setTime(cutoff).build());
// Seek to a snapshot.
subscriptionAdminClient.seek(SeekRequest.newBuilder()
.setSubscription(subName.toString()).setSnapshot(snapshotName.toString()).build());const sub = pubSubClient.subscription('orders-sub');
// Seek to a timestamp.
await sub.seek(cutoff);
// Seek to a snapshot.
await sub.seek('orders-snapshot');// Seek to a timestamp.
await subscriber.SeekAsync(new SeekRequest
{
SubscriptionAsSubscriptionName = subName,
Time = Timestamp.FromDateTime(cutoff),
});
// Seek to a snapshot.
await subscriber.SeekAsync(new SeekRequest
{
SubscriptionAsSubscriptionName = subName,
SnapshotAsSnapshotName = snapshotName,
});sub = pubsub_client.subscription "orders-sub"
# Seek to a timestamp.
sub.seek cutoff
# Seek to a snapshot.
sub.seek snapshotTimestamp clamping
Seeking before the retained window clamps — it is not an error. A Seek to a timestamp older
than the earliest retained message does not fail; it clamps to the earliest retained message and
replays from there. Because per-resource retention is itself clamped to the broker's
Store.MaxRetention, what is "retained" depends on the broker ceiling. Don't rely on a pre-window
seek returning an error to detect "too far back" — it silently starts at the oldest available
message. See Reliability.
Replay cap
A single Seek replays at most CONNECTORS_GCP_MAX_SEEK_REPLAY messages (default 1,000,000).
Hitting the replay cap stops at the cap and logs a WARN — there is no silent loss. You simply do
not replay beyond the limit in one seek. Raise CONNECTORS_GCP_MAX_SEEK_REPLAY, or seek in smaller
windows, if you need to replay more.
Snapshots
A snapshot captures a subscription's current cursor so you can seek back to it later without knowing an exact timestamp:
CreateSnapshot(subscription)records the cursor as a registry record.Seek(subscription, snapshot)rewinds to that captured cursor.- Snapshots have a 7-day default expiry and are swept hourly.
UpdateSnapshotmay change thelabelsandexpire_time.
You cannot snapshot a detached subscription. CreateSnapshot on a subscription whose topic has
been deleted or detached returns FAILED_PRECONDITION. Snapshot before you detach.
Related
Was this page helpful?
Schema Validation
Avro and Protobuf schema enforcement over KubeMQ — CreateSchema, topic schema_settings, enforce-on-publish, ≤300 KB definitions, and revisions.
Subscribing
Consume over KubeMQ — Pull vs StreamingPull with flow control, ack-deadline leases, ModifyAckDeadline nack/extend, and exactly-once with its node-local note.