Rebuild from the start¶
A pipeline is rebuilt by moving its streamlets' consumer groups back to the earliest offset of every topic they read, so they read their inputs again from the beginning. Kafka refuses to move a group that still has members, so the streamlets are stopped first, reset, and started again. The operator does the Kafka work, over the same connection each streamlet uses.
Each inlet has its own consumer group, <pipeline>.<streamlet>.<inlet>, so a reset can target some
streamlets and leave the rest where they are. Resetting reprocesses records; it does not delete what
downstream topics already hold, so downstream streamlets see those records again.
Stop the streamlets¶
Set replicas to 0 for every streamlet to reset, through deploy-time configuration, and apply the
regenerated resource:
# stop.conf
flow.streamlets.router { replicas = 0 }
flow generate blueprint.conf --descriptors flow --conf prod.conf --conf stop.conf \
--image router=registry.example.com/cart-router:1.2 -n shop | kubectl apply -f -
kubectl -n shop get pods -l flow.ankka.thinkmorestupidless.com/pipeline=cart
Wait until the streamlets' pods are gone. Scaling the Deployment directly does not work: the operator
restores the resource's replicas, and flow reset reads replicas from the resource.
Request the reset¶
flow reset cart -n shop
reset requested for 'cart': 1b7c2f0e-9a4d-4d3e-8f55-2f1c0b6e9a71
With no --streamlet, every streamlet with an inlet is reset. --streamlet router, repeatable, limits
it to the named ones. flow reset refuses, and changes nothing, when a target still has replicas other
than 0 or has pods left, when a named streamlet does not exist or has no inlets, or when the pipeline is
not found. Every refusal is listed in the CLI reference.
The request is an annotation on the AnkkaFlow. The operator checks the same guards again, because the
CLI's view of the pods may be stale: while a target runs it records ResetRefused with
reset <id> waits: … and tries again every few seconds.
Check that it happened¶
The operator records one event per consumer group:
kubectl -n shop get events --field-selector involvedObject.kind=AnkkaFlow,reason=ResetOffsets
REASON OBJECT MESSAGE
ResetOffsets ankkaflow/cart cart.router.in: 3 partition(s) of 'shop.cart-events.v1' to earliest
A group that could not be reset is ResetOffsetsFailed, with Kafka's reason; the pipeline itself is not
failed by it. When every group has been handled, the operator writes the request's id to the annotation
flow.ankka.thinkmorestupidless.com/reset-offsets-done, so the same request is never carried out again,
even after the operator restarts. A new reset is a new request with a new id.
Start the streamlets again¶
Apply the resource without the stop configuration:
flow generate blueprint.conf --descriptors flow --conf prod.conf \
--image router=registry.example.com/cart-router:1.2 -n shop | kubectl apply -f -
The streamlets start at the earliest offsets, and their lag, labelled
client_id="<pipeline>.<streamlet>.<inlet>", starts at the size of each input and falls to zero as they
catch up. See Observe a pipeline.