Sidecar¶
Image ankka-flow-sidecar, published as ghcr.io/thinkmorestupidless/ankka-flow-sidecar and built
locally by sbt sidecar/docker:publishLocal. It runs beside every streamlet's process and owns
everything Kafka: subscribing, batching, producing, committing after the write, consumer groups, lag and
stall warnings. In a cluster the operator adds it to each pod from its own FLOW_SIDECAR_IMAGE setting;
a pipeline never names it. On a laptop it runs in Docker Compose beside a process on the host. How it
behaves is explained in The sidecar.
Environment¶
Each variable can also be given as a system property, the name lower-cased with _ replaced by .
(flow.process.address).
| variable | default | meaning |
|---|---|---|
FLOW_PROCESS_ADDRESS |
127.0.0.1:9010 |
host:port of the streamlet's process. The operator sets 127.0.0.1:9010; a compose file sets host.docker.internal:9010 |
FLOW_CONFIG_DIR |
/etc/flow/config |
holds descriptor.json and streamlet.conf |
FLOW_STATE_DIR |
/tmp/flow |
where the ready and alive files are written |
FLOW_METRICS_PORT |
2050 |
the Prometheus metrics port |
FLOW_DISCOVERY_TIMEOUT |
60s |
the deadline of one discovery call; attempts continue forever |
FLOW_STALL_WARNING_AFTER |
5m |
how long a partition may go without committing before a stall warning |
FLOW_RECONNECT_MAX_BACKOFF |
30s |
the ceiling of the backoff between conversations after a failure |
FLOW_POD_NAME, FLOW_POD_NAMESPACE |
set by the operator | the pod a stall warning is recorded on |
KUBERNETES_SERVICE_HOST, KUBERNETES_SERVICE_PORT |
set by Kubernetes | with the pod name and namespace, stall warnings become Events; otherwise they are log lines |
Durations are HOCON durations: 500ms, 30s, 5m. The process container receives
FLOW_PROCESS_PORT (9010) and nothing else from the platform.
descriptor.json¶
The streamlet's descriptor as deployed: in a cluster the resource's descriptor, locally the file
uv run descriptor wrote. The sidecar compares it with what the process answers at discovery. Its format
is in the descriptor reference.
streamlet.conf¶
HOCON naming the pipeline, the streamlet, its resolved parameters, and one entry per port with that port's topic and its own Kafka connection. In a cluster the operator renders it into the streamlet's Secret. On a laptop you write it beside a compose file:
# The sidecar's configuration for the compose network. In a cluster the operator renders the same
# file from the pipeline resource. descriptor.json beside it is written by `uv run descriptor`.
flow {
pipeline = cart
streamlet = router
config = { review-threshold = 100 }
inlets {
in {
topic = "shop.cart-events.v1"
group = "cart.router.in"
client-id = "cart.router.in"
bootstrap.servers = "kafka:9092"
connection-config {}
consumer-config { auto.offset.reset = earliest }
batch { max-records = 100, max-bytes = 1 MiB }
}
}
outlets {
valid {
topic = "cart.valid-carts"
client-id = "cart.router.valid"
bootstrap.servers = "kafka:9092"
connection-config {}
producer-config {}
}
review {
topic = "cart.review-carts"
client-id = "cart.router.review"
bootstrap.servers = "kafka:9092"
connection-config {}
producer-config {}
}
}
}
| key | required | meaning |
|---|---|---|
flow.pipeline, flow.streamlet |
yes | the pipeline id and the streamlet's name |
flow.config |
no | parameter values; a parameter with a default may be left out. Sent to the process as Start.config_json |
inlets.<name>.topic, outlets.<name>.topic |
yes | the Kafka topic's name |
…bootstrap.servers |
yes | the brokers, per port |
inlets.<name>.group |
no | the consumer group; default <pipeline>.<streamlet>.<inlet> |
…client-id |
no | the Kafka client id; default <pipeline>.<streamlet>.<port> |
…connection-config |
no | client properties for this port: security.protocol, sasl.*, ssl.* |
inlets.<name>.consumer-config, outlets.<name>.producer-config |
no | consumer or producer properties |
inlets.<name>.batch.max-records, …max-bytes |
no | the largest batch sent to the process; default 100 and 1 MiB |
Every inlet reads with auto.offset.reset = earliest unless its consumer-config says otherwise, so a
new pipeline reads its inputs from the start. allow.auto.create.topics is always false: reading a
topic never creates it. The inlets and outlets must be exactly the descriptor's ports.
Start-up and exit codes¶
- Read the environment,
descriptor.jsonandstreamlet.conf, and check that the file's ports and parameters match the descriptor. Any problem is logged and the sidecar exits with code2. - Call
Discoveron the process until it answers, retrying from 500 ms, doubling to 10 s, forever. - Check the answer: the protocol major version must equal the sidecar's and its minor must not be later;
the descriptor must be valid and equal to the deployed one. On any problem it sends every problem to
the process's
ReportError, logs them, and exits with code1. - Open a
Runconversation, subscribe every inlet, and writereadyonce every inlet is subscribed and its topic exists.
On a failure after start-up the sidecar never exits: it discards the batches in flight, removes ready,
waits (500 ms, doubling to FLOW_RECONNECT_MAX_BACKOFF), repeats discovery, and resumes from the last
committed offsets. It exits 0 when it is stopped. The protocol itself is in
the protocol reference.
Probes¶
Rendered by the operator on the sidecar container only; the process container has no probes.
| probe | command | timing |
|---|---|---|
| readiness | test -f /tmp/flow/ready |
every 5 s |
| liveness | the age of /tmp/flow/alive is under 15 s |
every 10 s, after 20 s |
ready exists while a conversation runs and every inlet is subscribed to a topic that exists. It is
removed the moment the conversation fails or the sidecar is stopping. alive is touched every second.
Metrics¶
Prometheus text on :2050/metrics, from the JMX exporter agent. Every Kafka client's id is
<pipeline>.<streamlet>.<port>, so client_id attributes lag and throughput to one port of one
streamlet.
| metric | labels | meaning |
|---|---|---|
kafka_consumer_consumer_fetch_manager_metrics_records_lag |
client_id, topic, partition |
records behind, per inlet partition |
kafka_consumer_consumer_fetch_manager_metrics_records_lag_max |
client_id, topic, partition |
the maximum lag in the current window |
kafka_consumer_consumer_fetch_manager_metrics_records_consumed_rate |
client_id, topic |
records read per second |
kafka_producer_producer_metrics_record_send_rate |
client_id, topic |
records written per second, per outlet |
ankka_flow_sidecar_in_flight |
inlet, partition |
1 while a batch of that partition is with the process |
ankka_flow_sidecar_stalled_seconds |
inlet, partition |
how long the partition has gone without committing; 0 when nothing is outstanding |
Kafka reports topic names in these labels with dots replaced by underscores:
shop.cart-events.v1 appears as shop_cart-events_v1. The pod template carries
prometheus.io/scrape: "true" and prometheus.io/port: "2050".
Stall warnings¶
When ankka_flow_sidecar_stalled_seconds for a partition first passes FLOW_STALL_WARNING_AFTER, the
sidecar records one warning for that stall, reason PartitionStalled, with a note naming the pipeline,
streamlet, inlet and partition, how long it has not committed, and the last error. In a pod it is an
events.k8s.io/v1 Warning Event regarding the pod, posted with the pipeline's service account. Elsewhere,
or if posting fails, it is a warn log line. The stall clears when a batch of that partition commits.
Logs¶
One line per event on stdout at info: the start-up line naming the pipeline, streamlet and process
address, each discovery attempt and its result, each conversation's id, partitions assigned, revoked and
lost, each stream failure with its cause, each reconnect with its backoff, and readiness. A missing input
topic is logged as inlet '<name>': topic '<topic>' does not exist; not ready until it does. No record
value is ever logged.