Python SDK¶
The package is ankka-flow, imported as ankka_flow, in
sdks/python. It needs
Python 3.12 or later and depends on grpcio and protobuf. It is published to PyPI as
ankka-flow (uv add ankka-flow), versioned as the
ankka-flow release it belongs to. It is typed (py.typed) and checked with
mypy --strict.
Names¶
Everything below is importable from ankka_flow unless another module is named.
| name | what it is |
|---|---|
Streamlet |
the base class of every streamlet |
JsonInlet, JsonOutlet |
ports with a JSON contract |
StringParameter, IntegerParameter, DoubleParameter, BooleanParameter, DurationParameter, MemorySizeParameter |
typed configuration parameters |
Parameter |
the base class of the parameters |
Config |
a streamlet's resolved parameter values |
Record, Batch, Emit |
a record, a batch of records, a record for an outlet |
serve |
runs a streamlet for the sidecar |
json |
loads(bytes) and dumps(obj) -> bytes |
PROTOCOL_VERSION |
the protocol version the SDK speaks, "1.0" |
ankka_flow.testkit.Harness |
runs a streamlet over in-memory inlets and outlets |
ankka_flow.testkit.hash_partitioner |
a stable key-to-partition function for the Harness |
Streamlet¶
class Streamlet(ABC):
name: ClassVar[str]
description: ClassVar[str] = ""
config: Config
def process(self, batch: Batch) -> Iterable[Emit]: ...
nameis required: 1 to 63 of[a-z0-9-], not starting or ending with-.- Ports and parameters are class attributes, found when the subclass is defined, including those
inherited from a base class. A port name used twice, or a parameter key used twice, raises
TypeErrorat class definition. configholds the parameter values: the declared defaults until the sidecar'sStartapplies the deployed values.processis abstract. It is called once per batch on a worker thread; never twice at once for one partition, possibly concurrently for different partitions. Returning acknowledges the batch; raising fails it; yielding nothing for a record skips it. Yielding something other than anEmit, or an emit to an undeclared outlet, fails the batch.inlets(),outlets()andparameters()are class methods returning the declared ports and parameters.configure(values)resolves a mapping of parameter values intoconfig.
Ports¶
JsonInlet(name: str, *, schema_name: str)
JsonOutlet(name: str, *, schema_name: str)
name is the wire name, matching [a-z][a-z0-9-]{0,62}. schema_name is required and must not be
empty. Every port has format ("json"), schema_name and fingerprint, the Base64 of the SHA-256
of the schema name.
JsonOutlet.emit(record: Record | None = None, *, value: bytes | None = None,
key: bytes | None = <unset>, headers: list[tuple[str, bytes]] | None = None) -> Emit
With a record, emit forwards it: the same key, headers and value, and the same offset and timestamp,
with any keyword argument replacing its part. key=None removes the key. Without a record, value is
required and the other parts default to no key and no headers. Calling neither raises ValueError.
Parameters¶
StringParameter(key: str, *, default: str | None = None, description: str = "")
Every parameter class takes the same arguments. key matches [a-z][a-z0-9-]*. A parameter with no
default is required when the pipeline is deployed; reading it from config when it has no value
raises KeyError. A default that is not a value of the parameter's type raises ValueError at
declaration.
| class | descriptor type | config[param] returns |
accepts |
|---|---|---|---|
StringParameter |
STRING |
str |
any text |
IntegerParameter |
INTEGER |
int |
an integer |
DoubleParameter |
DOUBLE |
float |
a number |
BooleanParameter |
BOOLEAN |
bool |
true or false |
DurationParameter |
DURATION |
datetime.timedelta |
a HOCON duration: 100 ms, 5m, 1.5 s; a bare number is milliseconds |
MemorySizeParameter |
MEMORY_SIZE |
int, in bytes |
a HOCON size: 1 MiB, 512k, 10MB; a bare number is bytes |
A default may be given as the Python type (default=timedelta(seconds=5), default=True) or as text
(default="5s").
Config[param] returns the typed value; Config.get(key) returns a value by key, or None;
Config.as_dict() returns every resolved value.
Records¶
@dataclass(frozen=True)
class Record:
value: bytes
key: bytes | None = None
headers: list[tuple[str, bytes]] = []
offset: int = -1
timestamp_ms: int = 0
@dataclass(frozen=True)
class Batch: # iterable over its records; len() is their count
inlet: str
partition: int
records: list[Record]
@dataclass(frozen=True)
class Emit:
outlet: str
record: Record
A record is exactly what Kafka holds; nothing is decoded. Headers keep their order. A batch's records come from one partition of one inlet, in offset order. The SDK sends each emit to the sidecar as it is yielded; the sidecar holds a batch's emits until the batch is acknowledged and writes none of them if it fails.
serve¶
serve(streamlet: Streamlet, port: int | None = None, block: bool = True, workers: int = 32) -> Server
Serves the protocol's Discovery and Streamlet services on 127.0.0.1:port, where port defaults
to FLOW_PROCESS_PORT, or 9010 when that is unset. It binds loopback only and raises OSError when the
port cannot be bound. workers is the number of threads that run process.
Discoveranswers with the streamlet's descriptor.ReportErrorlogs each problem the sidecar reports atERRORlevel.- One conversation runs at a time. A new
Runends the previous one, because the sidecar has reconnected.Startapplies the deployed parameter values;Stoplets in-flight batches finish and ends the conversation. - Messages are limited to 16 MiB in each direction.
- With
block=True,SIGTERMandSIGINTstop the server andservereturns when it has stopped. Withblock=Falseit returns at once; the returnedServerhasport,stop(grace=2.0)andwait().
When the root logger has no handlers, serve configures logging at INFO. The SDK logs under the
logger name ankka_flow.
ankka_flow.testkit¶
Harness(streamlet: Streamlet, config: Mapping[str, object] | None = None)
| member | what it does |
|---|---|
inlet(name).put(*, value, key=None, headers=None) |
queues a record on a declared inlet |
run(partitions=None, max_records=None) |
processes everything queued: per inlet, per partition, batches in offset order |
outlet(name).records |
the records emitted to a declared outlet, in order |
skipped |
records of successful batches that no emit was derived from |
failures |
a Failure(batch, error) per failed batch; its emits are discarded |
batches |
every batch given to process, in order |
partitions maps a key to a partition number (default: everything on partition 0); max_records
caps a batch (default: one batch per partition). hash_partitioner(n) places keys on n partitions by
CRC-32 of the key, with a keyless record on partition 0. An undeclared inlet, outlet or configuration
key raises KeyError. See Test a streamlet.
Commands¶
The package installs two console scripts.
descriptor¶
uv run descriptor [--out flow/descriptor.json] [--check]
Loads the streamlet named module:Class by the FLOW_STREAMLET environment variable or, when that is
unset, by pyproject.toml, validates its descriptor and writes it as canonical JSON. src/ is put on
the import path when it exists. --check writes nothing and exits 1 when the file is missing or
differs. An invalid declaration prints each problem and exits 1.
[tool.ankka-flow]
streamlet = "cart_router.router:CartRouter"
The descriptor's sdk block is {"name": "ankka-flow-python", "version": <the SDK's version>}. See
Descriptor.
conformance¶
uv run conformance
Serves the SDK's conformance reference streamlet on FLOW_PROCESS_PORT (default 9010) and runs the
platform's conformance suite against it with sbt. It must run inside the ankka-flow repository, with
sbt and a JDK. ANKKA_FLOW_CONFORMANCE_ONLY runs only the cases whose names start with its value. See
Adding a language SDK.