Consuming for Analytics

Consuming BookNest Order Events for Analytics

BookNest's analytics consumer reads the order events with the new protocol, filters on the event-type header set by the shop's producer, keeps running figures, and commits after each batch: at-least-once, which is safe because recomputing a batch gives the same figures.

booknest_analytics.py: order metrics from the event streamPython
"""booknest_analytics.py: order metrics computed from the event stream, committed per batch."""
import json, time
from collections import Counter
from confluent_kafka import Consumer
c = Consumer({"bootstrap.servers": "localhost:33092", "group.id": "booknest-analytics",
              "group.protocol": "consumer", "auto.offset.reset": "earliest",
              "enable.auto.commit": False})
c.subscribe(["booknest.order-events"])
types, totals, cancelled, idle, start = Counter(), {}, set(), 0, time.time()
while idle < 5:                                   # five empty polls once data has flowed
    batch = c.consume(5000, timeout=1)
    idle = 0 if batch else idle + bool(types)
    for m in batch:
        kind = dict(m.headers() or []).get("event-type", b"").decode()
        if not kind:
            continue                              # test records without the header
        types[kind] += 1
        if kind == "order_placed":
            totals[m.key()] = json.loads(m.value())["total"]
        elif kind == "order_cancelled":
            cancelled.add(m.key())
    if batch:
        c.commit(asynchronous=False)              # at least once: after processing the batch
c.close()
gross = sum(total for order, total in totals.items() if order not in cancelled)
print(f"{sum(types.values()):,} events in {time.time() - start - 5:.1f} s:", dict(types))
print(f"orders {len(totals):,}, cancelled {len(cancelled):,}, revenue kept {gross:,.2f}")
Output
390,737 events in 4.4 s: {'order_placed': 100000, 'order_paid': 94106, 'order_cancelled':
  5890, ...}
orders 100,000, cancelled 5,890, revenue kept 3,270,268.02

The consumer read about 90,000 events a second and skipped the test records of Transactions and Exactly-Once, which carry no header. The revenue equals the jq 133,477 sum of total over JSON, Columnar and Binary Formats's orders.jsonl for orders not cancelled. More copies with the same group.id scale it out, each holding figures for its own partitions only; per-key state belongs in Kafka Streams 129 or Flink 129 (Kafka Streams and ksqlDB and Apache Flink).