Observe a pipeline¶
A pipeline reports at three levels: the AnkkaFlow's status says whether the whole pipeline runs, the
events on it say what the operator did and why, and each pod's sidecar says what one streamlet is doing
through its log and its metrics.
The phase¶
kubectl -n shop get aflow
kubectl -n shop get aflow cart -o wide -w
NAME PIPELINE PHASE DETAIL AGE
cart cart Ready 4m
| phase | meaning |
|---|---|
Pending |
a streamlet is still rolling out |
Ready |
every streamlet has all its pods ready and every topic exists |
Degraded |
an unmanaged topic is missing or Kafka is unreachable, or a streamlet has fewer ready pods than desired |
Failed |
the operator refused the resource and applied nothing; DETAIL lists why |
DETAIL joins the reasons with ;, such as topic 'shop.cart-events.v1' does not exist or
router: 1 of 3 ready. The full status, per streamlet and per topic:
kubectl -n shop get aflow cart -o jsonpath='{.status}' | jq .
A streamlet's pod is ready when its sidecar is: a conversation with the process is running and every inlet is subscribed to a topic that exists. The process container has no probes of its own.
Events¶
The operator records what it did on the AnkkaFlow:
kubectl -n shop get events --field-selector involvedObject.kind=AnkkaFlow
TopicCreated, StreamletRolled, StreamletRemoved and ResetOffsets are Normal. Every Warning names
something to act on; each is described in Troubleshooting. Stalled partitions are
recorded by the sidecar on its own pod:
kubectl -n shop get events --field-selector reason=PartitionStalled
Logs¶
Every pod has two containers: sidecar and process.
kubectl -n shop logs deploy/flow-cart-router -c sidecar
kubectl -n shop logs deploy/flow-cart-router -c process
The sidecar logs discovery, each conversation, partitions assigned and revoked, every stream failure with its cause, and every reconnect. When the sidecar refuses to start because the process does not match the deployed descriptor, the problems are in both logs: the sidecar sends them to the process before it exits. No record value is ever logged.
Lag and throughput¶
Each sidecar serves Prometheus metrics on port 2050. Every Kafka client's id is
<pipeline>.<streamlet>.<port>, so lag is attributed to one streamlet's inlet:
kubectl -n shop port-forward deploy/flow-cart-router 2050 &
curl -s localhost:2050/metrics | grep records_lag
kafka_consumer_consumer_fetch_manager_metrics_records_lag{client_id="cart.router.in",partition="0",topic="shop_cart-events_v1"} 0.0
Kafka writes topic names in these labels with dots replaced by underscores. The metrics worth watching:
| metric | watch for |
|---|---|
kafka_consumer_consumer_fetch_manager_metrics_records_lag |
lag per inlet partition that grows and does not fall |
kafka_consumer_consumer_fetch_manager_metrics_records_consumed_rate |
an inlet that stops reading |
kafka_producer_producer_metrics_record_send_rate |
emits per outlet |
ankka_flow_sidecar_stalled_seconds |
a partition that has not committed for a long time |
ankka_flow_sidecar_in_flight |
a batch that stays with the process |
Every metric is listed in the sidecar reference.
Scraping with Prometheus¶
Streamlet pods carry prometheus.io/scrape: "true" and prometheus.io/port: "2050", and the port is
named metrics. A Prometheus that discovers pods by those annotations scrapes every sidecar without
further configuration. With the Prometheus Operator, a PodMonitor selecting
app.kubernetes.io/managed-by: ankka-flow on port metrics does the same:
apiVersion: monitoring.coreos.com/v1
kind: PodMonitor
metadata:
name: ankka-flow
namespace: shop
spec:
selector:
matchLabels:
app.kubernetes.io/managed-by: ankka-flow
podMetricsEndpoints:
- port: metrics
Pods also carry flow.ankka.thinkmorestupidless.com/pipeline and
flow.ankka.thinkmorestupidless.com/streamlet labels, which a relabelling rule can turn into series
labels.
Stalled partitions¶
A batch the process fails is redelivered from the last commit, indefinitely; nothing is skipped. A batch
that fails every time therefore stalls its partition. It shows three ways: that partition's lag grows,
ankka_flow_sidecar_stalled_seconds rises, and after FLOW_STALL_WARNING_AFTER (five minutes by default)
the sidecar records one PartitionStalled Warning on its pod, naming the inlet, the partition and the
last error. The fix is in the streamlet's code: skip the record by acknowledging without emitting, or
correct what makes it fail. See Delivery and failure.