One of the more challenging aspects I encountered while working with stream processing was the task of performing near-merge or as-of joins on Kafka streams.
Especially when topics exist on multiple Kafka clusters, ruling out the option to connect and subscribe to multiple topics using a single consumer.
⟶ 
Some frameworks seem promising at first, mentioning “joining or merging” of event streams in the documentation. Or declaring support for arbitrary operations.
In my experience however, the following occurs:
As a consequence, the functionalities that support arbitrary stateful operations (like nearby joining of topics) were implemented in-house.
The reusable data-flow model was open-sourced first as Snapstream, then as the async successor Slipstream.
If you want to follow along, run a Kafka broker as shown in the previous post, and install Slipstream (Python 3.10+):
pip install slipstream-async
pip install slipstream-async[kafka] # Topic via aiokafka
pip install slipstream-async[cache] # Cache via rocksdict
The basic concepts are:
AsyncIterableCallablehandle binds sources and sinks to user-defined handler functionsstream processes each source in parallelHandlers may be sync or async. Start the flow with asyncio.run(stream()). A plain sync iter alone does not parallelize — wrap it so other coroutines can run:
from asyncio import sleep
async def async_iterable(it):
for msg in it:
await sleep(0.01)
yield msg
A Topic is both an AsyncIterable (consume) and a Callable (produce):

When the processing logic of your stream solely relies on the message itself, without requiring any additional data, it is referred to as “stateless processing”.
Any other type of processing usually involves caching and retrieving state. RocksDB is a popular state store that’s used in Spark, ksqlDB, Faust and other frameworks. Slipstream wraps it through rocksdict as Cache:
from slipstream import Cache
cache = Cache('state/db')
# value '🏆' is stored under the key 'prize'
cache('prize', '🏆')
Cache is also callable, so we can use it as our sink. Persisted data survives restarts.
At this point, we’re ready to perform some type of arbitrary operation on our data:
from datetime import datetime as dt
weather_messages = iter([
{'timestamp': dt(2023, 1, 1, 10), 'value': '🌞'},
{'timestamp': dt(2023, 1, 1, 12), 'value': '⛅'},
{'timestamp': dt(2023, 1, 1, 13), 'value': '🌧'},
])
activity_messages = iter([
{'timestamp': dt(2023, 1, 1, 10, 30), 'value': 'swimming'},
{'timestamp': dt(2023, 1, 1, 11, 30), 'value': 'walking home'},
{'timestamp': dt(2023, 1, 1, 12, 30), 'value': 'shopping'},
{'timestamp': dt(2023, 1, 1, 13, 10), 'value': 'lunch'},
])
Let’s try and figure out the weather at the time of each activity:
from asyncio import run, sleep
from slipstream import Cache, handle, stream
weather_cache = Cache('state/weather')
async def async_iterable(it):
for msg in it:
await sleep(0.01)
yield msg
@handle(async_iterable(weather_messages), sink=[weather_cache])
def handle_weather(w):
unix_ts = w['timestamp'].timestamp()
yield unix_ts, w
@handle(async_iterable(activity_messages), sink=[print])
def handle_activity(a):
unix_ts = a['timestamp'].timestamp()
for w in weather_cache.values(backwards=True, from_key=unix_ts):
yield a['value'], w['value']
break
yield a['value'], '?'
run(stream())
('swimming', '🌞')
('walking home', '🌞')
('shopping', '⛅')
('lunch', '🌧')
What we just did:
handle_weather functionhandle_activity functionNow you might be wondering two things…
The generators can easily be replaced with Topic instances (aiokafka under the hood; config keys are snake_case):
from slipstream import Topic
from slipstream.codecs import JsonCodec
weather_messages = Topic('weather', {
'bootstrap_servers': 'localhost:29091',
'auto_offset_reset': 'earliest',
'group_instance_id': 'weather',
'group_id': 'weather',
}, codec=JsonCodec())
Change the handler so that it reads the value from the Kafka message object:
@handle(weather_messages, sink=[weather_cache])
def handle_weather(msg):
val = msg.value # attribute, not a method
unix_ts = val['timestamp'].timestamp()
yield unix_ts, val
When Topic is used as a sink, yield a key and value — for example yield None, val.
When weather updates arrive later than activity events, a naive join enriches with stale data.

Imagine if the “Cloudy” and “Rainy” events had a bit of a delay. If we did nothing, we’d send out incorrect weather states for the “Shopping” and “Lunch” activities:

('swimming', '🌞')
('walking home', '🌞')
('shopping', '🌞') # incorrect
('lunch', '⛅') # incorrect
We could defer processing so that late events can be joined successfully, which has its downsides:

('swimming', '🌞')
('walking home', '🌞')
# defer 'shopping' event
('shopping', '⛅')
# defer 'lunch' event
('lunch', '🌧')
Or we could choose to send out revisions on incorrect results, which works great together with Kafka’s log compaction:

('swimming', '🌞')
('walking home', '🌞')
('shopping', '🌞') # incorrect
('shopping', '⛅')
('lunch', '⛅') # incorrect
('lunch', '🌧')
Slipstream ships first-class downtime detection for that revision path: Checkpoint plus Dependency. Heartbeat the upstream stream, check the pulse on the dependent stream, pause while the dependency is behind, then seek and reprocess when it catches up.
from datetime import timedelta
from typing import cast
from slipstream import Cache, Topic, handle, stream
from slipstream.checkpointing import Checkpoint, Dependency
from slipstream.codecs import JsonCodec
from slipstream.core import READ_FROM_END
activity = Topic('activity', {
'bootstrap_servers': 'localhost:29091',
'auto_offset_reset': 'earliest',
'group_instance_id': 'activity',
'group_id': 'activity',
}, codec=JsonCodec(), offset=READ_FROM_END)
weather_stream = async_iterable(weather_messages)
checkpoints_cache = Cache('state/checkpoints', target_table_size=10000)
weather_cache = Cache('state/weather')
async def downtime_callback(c: Checkpoint, d: Dependency) -> None:
print('\tThe stream is automatically paused.')
async def recovery_callback(c: Checkpoint, d: Dependency) -> None:
offsets = cast(dict[str, int], d.checkpoint_state)
print(f'\tDowntime resolved, seeking back to {offsets}.')
await activity.seek({int(p): o for p, o in offsets.items()})
checkpoint = Checkpoint(
'activity',
dependent=activity,
dependencies=[Dependency(
'weather_stream',
weather_stream,
downtime_threshold=timedelta(hours=1),
)],
downtime_callback=downtime_callback,
recovery_callback=recovery_callback,
cache=checkpoints_cache,
)
In the weather handler, call checkpoint.heartbeat with the event time. In the activity handler, call checkpoint.check_pulse and pass partition offsets so recovery can seek:
@handle(weather_stream, sink=[weather_cache, print])
async def handle_weather(w):
ts = w['timestamp']
await checkpoint.heartbeat(ts)
yield ts.timestamp(), w
@handle(activity, sink=[print])
async def handle_activity(msg):
a = msg.value
ts = dt.strptime(a['timestamp'], '%Y-%m-%d %H:%M:%S')
unix_ts = ts.timestamp()
if downtime := await checkpoint.check_pulse(ts, **{
str(msg.partition): msg.offset
}):
print(f'\tDowntime detected: {downtime}')
for w in weather_cache.values(backwards=True, from_key=unix_ts):
yield a['value'], w['value']
break
yield a['value'], '?'
A typical run looks like this: one faulty enrichment slips out before the pause, then recovery seeks the activity topic and reprocesses so the join corrects itself — the same revision idea as before, without a hand-rolled Queue.
Downstream consumers still need to handle corrections (compaction or deduplication by key). Stateful aggregations must not double-count when a message is replayed.
Full docs: slipstream.readthedocs.io.