Skip to main content

ankka-flow for Python

Write ankka-flow streamlets in Python. You declare ports and parameters and implement process; the sidecar the platform runs beside your container owns everything Kafka.

from collections.abc import Iterable
from ankka_flow import Batch, Emit, IntegerParameter, JsonInlet, JsonOutlet, Streamlet, json, serve

class CartRouter(Streamlet):
    name = "cart-router"
    inlet = JsonInlet("in", schema_name="cart-events.v1")
    valid = JsonOutlet("valid", schema_name="cart-events.v1")
    review = JsonOutlet("review", schema_name="cart-events.v1")
    threshold = IntegerParameter("review-threshold", default=100)

    def process(self, batch: Batch) -> Iterable[Emit]:
        for record in batch:
            event = json.loads(record.value)
            yield (self.review if event["total"] > self.config[self.threshold] else self.valid).emit(record)

if __name__ == "__main__":
    serve(CartRouter())          # 127.0.0.1:$FLOW_PROCESS_PORT (9010), loopback only
  • process runs once per batch on a worker thread. One batch per partition is in flight at a time, so it never runs twice for one partition at once; different partitions run concurrently.
  • Returning acknowledges the batch. Raising fails it, and the sidecar redelivers it from the last commit. Yielding nothing for a record skips it.
  • Records are bytes with a key and headers. The SDK decodes nothing; ankka_flow.json helps.

Descriptor

uv run descriptor writes flow/descriptor.json from the streamlet named by [tool.ankka-flow] streamlet = "module:Class" in your pyproject.toml (or $FLOW_STREAMLET). Commit it. uv run descriptor --check fails when it is stale. The format is proto/DESCRIPTOR.md.

Testing without Kafka

from ankka_flow.testkit import Harness

h = Harness(CartRouter(), config={"review-threshold": 50})
h.inlet("in").put(key=b"cart-1", value=b'{"total": 10}')
h.run()
assert [r.key for r in h.outlet("valid").records] == [b"cart-1"]

Developing the SDK

uv sync
uv run python scripts/proto.py   # copy ../../protocol into proto/ (committed) and generate src/ankka_flow/_proto/
uv run mypy                      # strict
uv run pytest -q                 # fixtures, the server against a scripted sidecar, the harness
uv run conformance               # the reference streamlet against the platform's conformance suite

proto/ must equal ../../protocol byte for byte; CI diffs it. tests/test_descriptor_fixtures.py proves this SDK writes every fixture descriptor exactly.

A new project starts from template/, which carries a compose file with Kafka and the sidecar for the laptop loop.

Conformance

uv run conformance serves the reference streamlet (ankka_flow._conformance) on 127.0.0.1:$FLOW_PROCESS_PORT and runs the sidecar's conformance suite against it from the repository root (sbt and a JDK are needed). Last run: every one of the 18 cases that apply to an SDK passes; the five violation.* and version.* cases are skipped, as they only run against the Scala double.

Two deliberate breaks prove the suite points at what broke:

ANKKA_FLOW_BREAK what it breaks cases that fail
keyless-empty-key a keyless emit carries an empty key exactly run.unkeyed-emit (SC-005)
ack-first the ack is sent before the batch's emits the 12 cases that emit, run.emits-precede-ack among them

Release files for ankka-flow 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for ankka-flow 0.1.0
File Size Uploaded
ankka_flow-0.1.0.tar.gz 68.9 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for ankka-flow 0.1.0
File Interpreter ABI Platform
ankka_flow-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 99.3 kB

Release files / ankka_flow-0.1.0.tar.gz

Download URL ankka_flow-0.1.0.tar.gz
Size 68.9 kB
Tags Source
SHA-256 checksum
How to use checksums
e63134e2f1650e0858a1ffaa0975aff9822aea864b286098c187d1bd16e107e1
BLAKE2b-256 checksum
How to use checksums
14e3a2c9e78a8f975a44a388a66de53bb23754d65dc680c6eb35af4199b7860f
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 28, 2026.

Transparency log

Release files / ankka_flow-0.1.0-py3-none-any.whl

Download URL ankka_flow-0.1.0-py3-none-any.whl
Size 30.4 kB
Tags Python 3
SHA-256 checksum
How to use checksums
b1e273e194788ae833efce82b2210f85477b3424a127954d8fc23165afa2d54d
BLAKE2b-256 checksum
How to use checksums
45f407dfacdc14eedf3a614b8c796317a451d197edec616f59979d249bd6df2e
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 28, 2026.

Transparency log

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page