Skip to content

Your first streamlet

The sample in samples/cart-router is a Python streamlet with one inlet of cart events, keyed by cart id, and two outlets. It sends each event to the review outlet when the cart's total is above a threshold, and to the valid outlet otherwise. This tutorial runs it on a laptop: Kafka and the sidecar in Docker, the streamlet itself as an ordinary Python process on the host.

You need Docker, sbt and uv; Build the tools lists them.

The streamlet

The streamlet declares its ports and parameter as class attributes and implements process, which receives one batch of records and yields emits:

from collections.abc import Iterable

from ankka_flow import Batch, Emit, IntegerParameter, JsonInlet, JsonOutlet, Streamlet, json


class CartRouter(Streamlet):
    name = "cart-router"
    description = "Routes cart events to the valid or review outlet."
    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,
        description="Carts with a total above this go to the review outlet.",
    )

    def process(self, batch: Batch) -> Iterable[Emit]:
        limit = self.config[self.threshold]
        for record in batch:
            event = json.loads(record.value)  # the SDK decodes nothing; this is the router's choice
            outlet = self.review if event["total"] > limit else self.valid
            yield outlet.emit(record)  # same key, same headers, same bytes

The process's entry point serves it on 127.0.0.1:$FLOW_PROCESS_PORT, where the sidecar finds it:

from ankka_flow import serve

from .router import CartRouter

if __name__ == "__main__":
    serve(CartRouter())

Test it without Kafka

Build the sidecar image once, from the repository root, then work in the sample's directory:

sbt sidecar/docker:publishLocal
cd samples/cart-router
uv sync
uv run pytest -q

The tests use the SDK's harness, which calls process with batches it builds and applies the protocol's rules. No Kafka, sidecar or network is involved. Test a streamlet describes the harness.

Check the descriptor

uv run descriptor --check

The SDK writes flow/descriptor.json from the streamlet's declaration. The descriptor is what a blueprint is checked against, and what the sidecar compares with the running process before it sends a single record. --check exits 1 when the committed file differs from the declaration; uv run descriptor rewrites it.

Start Kafka and the sidecar

docker compose up -d

The compose file starts a single-node Kafka, reachable from the host on localhost:9094, and the sidecar. The sidecar reads its configuration from flow/: the descriptor, and a streamlet.conf naming the pipeline, the streamlet, the topic behind each port and the Kafka address. In a cluster the operator writes that file; on a laptop it is committed beside the compose file.

The sidecar looks for the streamlet on host.docker.internal:9010 and keeps asking until it answers, so the two can start in either order.

Run the streamlet and send it events

uv run python -m cart_router.main &
uv run python produce.py

produce.py creates the three topics when they are missing (the compose Kafka does not create topics by itself, and there is no operator on a laptop), then writes fifty CloudEvents over ten cart ids to shop.cart-events.v1 with a plain Kafka client.

The sidecar subscribes to the input topic, sends batches to the router, writes the router's emits to cart.valid-carts and cart.review-carts, and commits the input offsets only once the broker has confirmed every write.

Kill it mid-stream

kill %1; sleep 1; uv run python -m cart_router.main &
uv run python verify.py

When the router goes away, the sidecar discards the batches in flight, commits nothing for them, and waits for the process to come back. When it does, the sidecar describes it again, opens a new conversation and resumes from the last committed offsets.

verify.py reads the input and both outlets from the beginning and exits 0 when every event is on the outlet its total chose, each cart's events are in the order they were produced, and the CloudEvents headers arrived intact. It reports repeats separately: delivery is at least once, so an event whose emits were written but whose offsets were not yet committed when the router died is delivered again. A repeat never reorders a cart. Delivery and failure explains why.

Clean up

kill %1
docker compose down

Where to go from here