Operator¶
The operator is one Deployment, ankka-flow-operator in namespace ankka-flow, image
ankka-flow-operator. It watches AnkkaFlow resources in every namespace, creates the managed Kafka
topics, and renders each streamlet as a Secret and a Deployment. It adds to a pipeline only what only it
can know: the sidecar image, which is its own setting, and the Kafka cluster settings, which it reads
from Secrets. Installing it is covered in Install the platform.
Settings¶
Each setting is read from a system property, then an environment variable, then its default. An empty value counts as unset.
| environment variable | system property | default | meaning |
|---|---|---|---|
FLOW_SIDECAR_IMAGE |
flow.operator.sidecar-image |
none | the sidecar image added to every streamlet pod. Unset, every pipeline is refused with SidecarImageMissing |
FLOW_KAFKA_CLUSTERS_NAMESPACE |
flow.operator.kafka-clusters-namespace |
FLOW_OPERATOR_NAMESPACE, else ankka-flow |
where the kafka-cluster-<name> Secrets are read from |
FLOW_OPERATOR_RESYNC_SECONDS |
flow.operator.resync-seconds |
300 |
how often every pipeline is reconciled even when nothing changed |
FLOW_OPERATOR_RETRY_MIN_BACKOFF_SECONDS |
flow.operator.retry-min-backoff-seconds |
1 |
the first retry delay after a failed reconcile; it doubles per attempt |
FLOW_OPERATOR_RETRY_MAX_BACKOFF_SECONDS |
flow.operator.retry-max-backoff-seconds |
300 |
the retry delay's ceiling |
FLOW_OPERATOR_MAX_CONCURRENT_RECONCILES |
flow.operator.max-concurrent-reconciles |
4 |
pipelines reconciled at once |
HOSTNAME, set by Kubernetes to the pod's name, is the reportingInstance of every event it records.
Upgrading the platform changes FLOW_SIDECAR_IMAGE. Every streamlet then rolls onto the new sidecar
image at its next reconcile, because the image is part of each pod template. A pipeline never names the
sidecar image.
Installed objects¶
The manifest at kustomization/components/operator/operator.yaml installs:
- the Namespace
ankka-flow; - the ServiceAccount
ankka-flow-operator; - a ClusterRole and ClusterRoleBinding
ankka-flow-operator; - the Deployment
ankka-flow-operator: one replica,Recreatestrategy,FLOW_SIDECAR_IMAGEandFLOW_KAFKA_CLUSTERS_NAMESPACEset, requesting 100m CPU and 256Mi memory.
The CRD is installed separately, from kustomization/components/crd/ankkaflow.yaml.
Permissions¶
The ClusterRole grants:
| API group | resources | verbs |
|---|---|---|
flow.ankka.thinkmorestupidless.com |
ankkaflows |
get, list, watch, patch, update |
flow.ankka.thinkmorestupidless.com |
ankkaflows/status |
get, update, patch |
apps |
deployments |
get, list, watch, create, update, patch, delete |
| core | secrets, serviceaccounts |
get, list, watch, create, update, patch, delete |
| core | pods |
get, list, watch |
rbac.authorization.k8s.io |
roles, rolebindings |
get, list, watch, create, update, patch, delete |
events.k8s.io |
events |
create, patch |
It holds create on events because it grants that permission to each pipeline's sidecars, and
Kubernetes forbids granting what the grantor does not hold. Every object it writes is applied
server-side with the field manager ankka-flow-operator, forcing conflicts: the resource is the source of
truth for every field the operator owns.
What starts a reconcile¶
The operator watches AnkkaFlow resources in every namespace, and Deployments labelled
app.kubernetes.io/managed-by: ankka-flow. Any change to either queues its pipeline, and every pipeline
is reconciled again every FLOW_OPERATOR_RESYNC_SECONDS. A reconcile that throws is retried with
backoff. A reset waiting for its targets' pods to go is looked at again every 5 s.
One reconcile¶
Rendering is pure: from the resource, the Kafka cluster Secrets, the observed Deployments, pods and topics, it produces a list of actions, which two executors carry out, one for Kubernetes and one for Kafka. In order:
- Refusals. The resource is refused, its phase set to
Failedwith every problem indetail, one Warning event per problem, and nothing else applied, when:FLOW_SIDECAR_IMAGEis unset (reasonSidecarImageMissing);- a topic names a cluster with no
kafka-cluster-<name>Secret, or the Secret has nobootstrap.servers(reasonRefused, as are all that follow); - a topic has no brokers after resolution;
- a managed topic has no
partitionsor noreplicasafter resolution; - a streamlet has no descriptor, or its descriptor does not parse or fails validation;
- a declared inlet is not bound, or a binding names a port the descriptor does not declare;
- a port names a topic id the resource does not declare.
- Topics. Each managed topic that does not exist is created with its partitions, replication and
topicConfig(TopicCreated). An existing managed topic is never altered: other partitions or replication isTopicDiffers, a differingtopicConfigentry isTopicSettingsIgnored. An unmanaged topic is only described; one that does not exist isTopicMissing. Topics are handled before any Deployment. - Per pipeline, the ServiceAccount, Role and RoleBinding
flow-<pipeline>. - Per streamlet, the Secret and the Deployment
flow-<pipeline>-<streamlet>, withStreamletRolledwhen the pod template's image or configuration hash changed. A labelled Deployment the spec no longer names is deleted with its Secret (StreamletRemoved). - A pending reset request, when the request annotation's id differs from the done annotation's.
While any target has
replicasother than 0 or pods left, it recordsResetRefusedand waits. A request naming an unknown streamlet recordsResetRefusedand is marked done. Otherwise it moves each target inlet's consumer group to the earliest offsets over that topic's own connection (ResetOffsets, orResetOffsetsFailedwith Kafka's reason, which is never a pipeline failure), then writes the done annotation. - Status.
phase,detail, and per streamlet and topic status; written only when it changed.
The phases and every event are listed in the resource reference. The operator logs each
event it records, at warn for a Warning and info otherwise.
Deletion¶
Deleting an AnkkaFlow deletes everything the operator rendered for it, through owner references. With
onDelete.managedTopics: Delete, the operator holds the finalizer
flow.ankka.thinkmorestupidless.com/managed-topics until it has deleted the managed topics. Unmanaged
topics are never deleted.