Writing a consumer¶
A consumer is any code that reads engine.snapshot(): a display, a monitor, a logger, a
second algorithm that works on the accumulated state. It needs no registration: it calls
snapshot() when it wants the latest state.
import threading
import time
import numpy as np
import frames2py
WIDTH, HEIGHT, BATCH, BATCHES = 640, 480, 10_000, 200
engine = frames2py.Engine((WIDTH, HEIGHT), "event_count", snapshot_interval_ms=1.0)
done = threading.Event()
seen = []
def produce():
rng = np.random.default_rng(seed=1)
for b in range(BATCHES):
events = np.zeros(BATCH, dtype=frames2py.EVENT_DTYPE)
events["t"] = np.arange(b * BATCH, (b + 1) * BATCH)
events["x"] = rng.integers(0, WIDTH, BATCH)
events["y"] = rng.integers(0, HEIGHT, BATCH)
engine.ingest(events) # never waits for the consumer
engine.stop() # publishes the events accumulated since the last publication
done.set()
def consume():
while True:
finished = done.is_set()
snapshot = engine.snapshot() # never waits for the producer
if snapshot is not None and (not seen or snapshot.meta.sequence != seen[-1]):
seen.append(snapshot.meta.sequence)
if finished:
break
time.sleep(0.001) # read at about the publication cadence
producer = threading.Thread(target=produce)
consumer = threading.Thread(target=consume)
consumer.start()
producer.start()
producer.join()
consumer.join()
final = engine.snapshot()
print("consumer saw new publications:", len(seen) > 0)
print("sequences only increase:", all(a < b for a, b in zip(seen, seen[1:])))
print("final watermark:", final.meta.watermark)
print("events ingested:", engine.stats.events_ingested)
consumer saw new publications: True
sequences only increase: True
final watermark: 1999999
events ingested: 2000000
The pattern¶
- Run the producer on its own thread and let it call
ingest()in its loop. It owns the Engine's ingest path; consumers never slow it down by waiting. - Poll at the publication cadence. Sleep about
snapshot_interval_msbetween reads. Reading faster only returns the same snapshot again. - Detect new publications by
meta.sequence. It increases by one per publication, for the Engine's lifetime. A jump of more than one means you skipped publications, which is normal for a consumer slower than the cadence. - Handle
None.snapshot()isNonebefore the first publication and afterreset()until the next. - Copy before you modify.
snapshot.frameis shared with every other consumer and read-only. Usesnapshot.copy()(orcopy(out=...)into a buffer you keep) before writing to the data or handing it to a library that ignores NumPy's read-only flag, such astorch.from_numpy.
What a consumer gets, and doesn't¶
- The latest state, not every event. Snapshots skip publications a slow consumer
missed, and a windowed kernel's frame covers only its own window. A consumer that needs
every event needs the events: call it from the producer's loop, next to
ingest()(as the recorder is used), accepting that its cost is then the producer's. - A consistent snapshot. Frame and metadata are from one publication, and the frame is complete.
- No ordering promise across consumers. Two consumers reading at the same moment may see different publications if one was published in between.
Keeping up¶
A consumer's own speed is its own concern. Frames2Py keeps only the latest snapshot and
never queues publications for a slow consumer, so a consumer that falls behind loses
publications, not memory. The flip side: holding references to old snapshots keeps their
frames alive; each publication is a fresh array (about 3.5 MiB for event_count at
1280x720).
On standard CPython, a consumer doing heavy pure-Python work competes with the producer for the GIL. Large NumPy operations can let other threads run (measured for frame copies on CPython 3.11), and a free-threaded 3.14t build has no GIL to compete for. See Lifecycle and threads.