Source and sink
Connect ubik to a Kafka topic as source and sink, replay or follow, per-partition offsets, and the transactional changelog a windowed query writes back.
Ubik speaks the Kafka protocol as both source and sink. It runs against the
Kafka, Redpanda, WarpStream or MSK you already have. A source is
--from kafka://host:port/topic; a sink is --to kafka://host:port/topic.
kafka://host:port/topic
kafka://host1:9092,host2:9092,host3:9092/topicThe port defaults to 9092. Several bootstrap brokers are comma-separated;
ubik discovers the rest of the cluster from any one it reaches.
A broker to talk to
Every command on this page runs against a local Redpanda. Any Kafka works the same; this is the shortest way to get one.
podman run -d --name redpanda -p 9092:9092 -p 8081:8081 \
docker.redpanda.com/redpandadata/redpanda:latest \
redpanda start --mode dev-container \
--kafka-addr PLAINTEXT://0.0.0.0:9092 \
--advertise-kafka-addr PLAINTEXT://localhost:9092 \
--schema-registry-addr 0.0.0.0:8081rpk, Redpanda's CLI, talks to it on localhost:9092. Create a topic with
three partitions and produce fifteen orders across three one-minute windows,
keyed by merchant so the partitioning is deterministic.
rpk topic create orders -p 3
python3 - <<'PY' | rpk topic produce orders --format '%k\t%v\n'
for w in range(3):
ts = "2026-07-18 10:%02d:00" % w
for merchant, amount, count in [("acme", 10, 3), ("globex", 3, 2)]:
for _ in range(count):
print('%s\t{"merchant":"%s","amount":%d,"event_time":"%s"}'
% (merchant, merchant, amount, ts))
PYThe two keys land on two of the three partitions; the third stays empty.
rpk topic describe orders -pPARTITION LOG-START-OFFSET HIGH-WATERMARK
0 0 9
1 0 6
2 0 0Read a topic
Point --from at the topic. With no --follow, ubik replays the topic from the
start to its current head and stops, closing every window at the end.
SQL="SELECT merchant,
TUMBLE(event_time, INTERVAL '1 MINUTE') AS w,
count(*) AS n,
sum(amount) AS total
FROM orders
GROUP BY merchant, TUMBLE(event_time, INTERVAL '1 MINUTE')"
ubik --from kafka://localhost:9092/orders "$SQL"The result is on stdout; status and stats are on stderr.
{"status":{"state":"done"}}
{"stats":{"rows_in":15,"rows_out":6,"late_dropped":0,"windows_closed":3,"rows_folded":15}}
{"merchant":"acme","w":"2026-07-18 10:00:00","n":3,"total":30}
{"merchant":"acme","w":"2026-07-18 10:01:00","n":3,"total":30}
{"merchant":"acme","w":"2026-07-18 10:02:00","n":3,"total":30}
{"merchant":"globex","w":"2026-07-18 10:00:00","n":2,"total":6}
{"merchant":"globex","w":"2026-07-18 10:01:00","n":2,"total":6}
{"merchant":"globex","w":"2026-07-18 10:02:00","n":2,"total":6}The table name in the SQL is the topic name. A JSON payload is self-describing, so no schema flag is needed; other wire formats are on the formats page.
Replay or follow
--follow tails the topic live instead of replaying to the head and stopping.
The difference is when a window closes.
| Mode | Window closes when |
|---|---|
| Replay (default) | Every remaining window closes at the end of the topic. |
--follow | A window closes only as the event-time watermark passes its end. |
The watermark is the minimum event time across the topic's partitions, so a single partition with no fresh events holds every window open. That is the correct behavior for exactly-once over a live stream, and the reason a follow can sit on a window a replay would have emitted.
Two flags bound it. --watermark-delay holds a window open past its end for
late rows within a grace period. --idle-partition-timeout stops a partition
from holding the watermark once the broker proves it is at the end of its log
and it has stayed quiet that long, so one permanently idle partition (the empty
partition 2 above, in a live topic) does not freeze every window. A partition
that still holds unread records is never excluded.
Write the changelog back
--to kafka://host:port/topic publishes the aggregation as a changelog. Each
row is keyed by its (group, window) tuple and upserted, so the topic's
latest-per-key state is the current aggregate.
rpk topic create orders_agg
ubik --from kafka://localhost:9092/orders "$SQL" \
--to kafka://localhost:9092/orders_aggRead it back. The key is the group tuple, the value is the row.
rpk topic consume orders_agg --offset start --num 6 --format '%k => %v\n'["acme","2026-07-18 10:00:00"] => {"merchant":"acme","w":"2026-07-18 10:00:00","n":3,"total":30}
["globex","2026-07-18 10:00:00"] => {"merchant":"globex","w":"2026-07-18 10:00:00","n":2,"total":6}
["acme","2026-07-18 10:01:00"] => {"merchant":"acme","w":"2026-07-18 10:01:00","n":3,"total":30}
["globex","2026-07-18 10:01:00"] => {"merchant":"globex","w":"2026-07-18 10:01:00","n":2,"total":6}
["acme","2026-07-18 10:02:00"] => {"merchant":"acme","w":"2026-07-18 10:02:00","n":3,"total":30}
["globex","2026-07-18 10:02:00"] => {"merchant":"globex","w":"2026-07-18 10:02:00","n":2,"total":6}Each closed window is wrapped in one Kafka transaction, so a read_committed
consumer sees a window atomically. The same changelog is also streamed to
stdout, so you can watch it while it publishes. How this survives a crash, and
why a re-emitted window upserts byte-identical rows over the first, is on the
exactly-once page.
Flags
| Flag | Meaning |
|---|---|
--from kafka://host:port/topic | Source topic. Comma-separate several bootstrap brokers. |
--to kafka://host:port/topic | Sink for the changelog. Defaults to stdout. |
--follow | Tail the topic live instead of replaying to the head and stopping. |
--watermark-delay <interval> | Hold a window open until end + interval for late rows. |
--idle-partition-timeout <interval> | Drop a proven-idle partition from the watermark after this long. |
--broker-timeout <interval> | Fail a run with BROKER_LOST after this much mid-run broker unreachability. Default 5 minutes. |
Authentication (TLS, SASL) is on the security page.
Exactly-once
What a checkpoint holds, what a kill -9 guarantees, and a run you can reproduce that stops mid-window and resumes with no row lost and none double-counted.
TLS and SASL
Authenticate to a Kafka broker with TLS and SASL through UBIK_KAFKA_ environment variables, which map one to one onto librdkafka properties.