Buffer ticks for downstream processing
Chapter 38 introduced pushing ticks to a queue instead of processing inline. This chapter builds that into a proper producer-consumer pipeline, and covers what to do when the consumer falls behind.
import queue, threading, time
class TickBuffer:
def __init__(self, maxsize: int = 10000):
self.q = queue.Queue(maxsize=maxsize)
self.dropped = 0
def put(self, tick: dict):
try:
self.q.put_nowait(tick)
except queue.Full:
self.dropped += 1 # never block the WebSocket thread
def consume_forever(self, handler):
while True:
tick = self.q.get()
try:
handler(tick)
except Exception:
logging.exception("Error processing tick")
finally:
self.q.task_done()
buffer = TickBuffer()
kws.on_ticks = lambda ws, ticks: [buffer.put(t) for t in ticks]
consumer_thread = threading.Thread(target=buffer.consume_forever, args=(process_tick,), daemon=True)
consumer_thread.start()
Why bounded queue + drop-on-full, not unbounded
An unbounded queue under sustained tick-processing lag just grows memory until the process dies — worse than losing some intermediate ticks. A bounded queue with put_nowait + a dropped-tick counter fails gracefully: you lose some ticks (rarely the last price matters as much as the most recent one, which you'll get next), and you get a visible signal (buffer.dropped climbing) that your consumer can't keep up, which is itself useful operational information.
Alert on sustained drops
def check_buffer_health(buffer: TickBuffer, threshold: int = 50):
if buffer.dropped > threshold:
send_alert(f"Tick buffer dropped {buffer.dropped} ticks — consumer is falling behind")
Never process ticks synchronously inside strategy logic that also places orders
Keep the tick-consumption thread separate from anything that blocks on network I/O (order placement, chapter 44+) for more than a few milliseconds — otherwise a slow order-placement call delays processing of every subsequent tick behind it in the queue.