Skip to content

CLI

flow verifies a blueprint against streamlet descriptors, writes the AnkkaFlow resource the operator runs, and asks the operator to reset a pipeline's consumer groups. It needs a JVM. Nothing installs it: from a clone of the repository, just cli (or sbt cli/stage) builds it at cli/target/universal/stage/bin/flow, and that directory goes on your PATH:

just cli
export PATH="$PWD/cli/target/universal/stage/bin:$PATH"
exit code meaning
0 success, or --help
1 refused or failed; every problem is on stderr, one per line
2 usage error: an unknown command or option, or a missing argument

verify and generate need no image, no network and no language runtime. reset needs a kubeconfig.

flow verify

flow verify <blueprint.conf> --descriptors <dir> [--conf <file>]...
option meaning
<blueprint.conf> the blueprint, HOCON; see the blueprint reference
--descriptors <dir> a directory whose *.json files are streamlet descriptors; other files are ignored
--conf <file> deploy-time configuration, HOCON; repeatable, later files win; see Configure at deploy time

It reads and validates every descriptor, parses the blueprint, checks it against the descriptors, and checks the --conf files against the result. On success it prints verified: <n> streamlets, <m> topics on stdout. Otherwise it refuses, listing every problem it found in one pass:

problem message
the directory is missing or empty descriptors: '<dir>' is not a directory, descriptors: '<dir>' holds no *.json files
a descriptor does not parse or fails validation (a bad name, a port declared twice, a fingerprint that does not match its schema name, an unsupported protocol version) descriptor <file>: <problem>
two descriptors declare one streamlet 2 descriptors declare streamlet '<name>'
the blueprint is not a file, or not HOCON blueprint: '<path>' is not a file, The blueprint file has an invalid format: …
a streamlet names a descriptor nobody declares Streamlet '<name>' names descriptor '<descriptor>', which no descriptor declares.
two streamlets share a name, or a name is illegal Duplicate streamlet names detected: …, Invalid streamlet name '<name>'. …
a port path names nothing '<path>' does not point to a known streamlet inlet or outlet, please try <suggestion>.
a producer is not an outlet, or a consumer is not an inlet '<path>' is not a valid producer for topic '<id>', must be an outlet. (and the consumer form)
an outlet and an inlet on one topic carry different contracts '<outlet>' (<contract>) is not compatible with '<inlet>' (<contract>).
a port uses a format other than json '<path>' uses format '<format>', which this version does not support; the only contract format is json.
an inlet is connected to nothing Inlet <streamlet>.<port> is not connected.
a port is bound to more than one topic '<path>' is bound to more than one topic: <ids>.
an unmanaged topic has producers Topic '<id>' is not managed but has producers <paths>; the platform only reads topics it does not own.
an unmanaged topic names no brokers and no cluster Topic '<id>' is not managed and names no bootstrap.servers or cluster.
an illegal topic or Kafka cluster name '<name>' is not a valid topic name, …, Invalid Kafka cluster name '<name>'. …
--conf does not parse --conf <file>: <reason>
--conf names a topic or streamlet the blueprint does not overrides name unknown topic '<id>', overrides name unknown streamlet '<name>'
--conf sets a parameter the descriptor does not declare streamlet '<name>': parameter '<key>' is not declared
a parameter has no default and no value streamlet '<name>': parameter '<key>' has no default and no value
a parameter's value is not its type streamlet '<name>': parameter '<key>' = <value> is not a <type>
replicas is not a number, or is negative streamlet '<name>': replicas must be a number, … must not be negative

An outlet connected to nothing is allowed. It is printed on stderr as a note (note: Outlet <streamlet>.<port> is not connected.) and does not change the exit code.

flow generate

flow generate <blueprint.conf> --descriptors <dir> [--conf <file>]...
              [--images <file>] [--image <name>=<ref>]...
              [--pipeline <id>] [--version <v>] [-n|--namespace <ns>] [-o|--output <file>]

Everything verify does, then it writes the AnkkaFlow resource as YAML, to stdout or to --output (which prints wrote <file> on stderr). Deploy-time configuration is merged over the blueprint here, so the resource says exactly what will run; the operator adds only the sidecar image and Kafka cluster settings.

option meaning
--images <file> a HOCON map of streamlet name to image reference
--image <name>=<ref> one streamlet's image; repeatable; wins over --images
--pipeline <id> the pipeline id; default blueprint.name, else the blueprint's file name up to its first .
--version <v> spec.version; default git describe --tags --always --dirty in the blueprint's directory, else unversioned
-n, --namespace <ns> metadata.namespace; without it the resource has none and kubectl uses its current namespace
-o, --output <file> write here instead of stdout

The resource's name and spec.pipeline are both the pipeline id. On top of verify's problems it refuses when:

problem message
a streamlet has no image Streamlet '<name>' has no image.
an --images file does not parse images: <reason>
an --image is not name=ref --image '<value>' is not name=reference
the pipeline id is not a DNS label of at most 40 characters pipeline id '<id>' must be 1-40 of [a-z0-9-], not starting or ending with '-'

An images file:

router = "ghcr.io/example/cart-router:0.3.1"
sink   = "ghcr.io/example/cart-sink:0.3.1"

A typical deployment:

flow generate blueprint.conf --descriptors flow --conf prod.conf \
  --image router=registry.example.com/cart-router:1.2 -n shop | kubectl apply -f -

spec.onDelete.managedTopics is always written as Keep. See the resource reference for every field generate writes.

flow reset

flow reset <pipeline> [--streamlet <name>]... [-n|--namespace <ns>]

Records a request that the named streamlets reread their inputs from the earliest offset. The operator carries it out; see Rebuild from the start.

option meaning
<pipeline> the AnkkaFlow's name
--streamlet <name> a streamlet to reset; repeatable; default every streamlet with an inlet
-n, --namespace <ns> the pipeline's namespace; default the kubeconfig's current namespace

It reads the AnkkaFlow and the pods labelled with its pipeline, then refuses when:

problem message
Kubernetes cannot be reached cannot reach Kubernetes: <reason>
no such pipeline no pipeline '<pipeline>' in namespace '<ns>'
a named streamlet does not exist, or has no inlets cannot reset offsets: no streamlet [<name>], … streamlet [<name>] has no inlets, so no consumer groups to reset
no streamlet has an inlet pipeline <pipeline> has no streamlets with inlets to reset
a target's replicas is not 0, or it still has pods cannot reset offsets while streamlets are running: [<name>] is not scaled to 0; [<name>] still has <n> pod(s). Stop them first: …

Otherwise it writes the annotation flow.ankka.thinkmorestupidless.com/reset-offsets with {"id":"<uuid>","streamlets":[…]} and prints reset requested for '<pipeline>': <uuid>. An empty streamlets list means every streamlet with an inlet.

flow version

flow version

Prints the CLI's version and the protocol version its resources carry, as flow <version>, protocol <major>.<minor>.