GenStage Aggregation

Copy Markdown View Source
Mix.install([
  {:ex_data_sketch, "~> 0.10"},
  {:gen_stage, "~> 1.0"}
])

Introduction

GenStage provides back-pressure-aware data exchange between producers and consumers. ExDataSketch.GenStage provides SketchConsumer, SketchProducer, and SketchStage for integrating sketches into GenStage pipelines.

Section 1: SketchConsumer -- Accumulate Events

SketchConsumer is a GenStage consumer that accumulates events into a sketch. It supports update, merge, estimate, get, and flush operations.

# Start a consumer that builds an HLL sketch
{:ok, consumer} = ExDataSketch.GenStage.SketchConsumer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14],
  subscribe_to: []
)

# Merge partial sketches into the consumer
partial = ExDataSketch.HLL.from_enumerable(1..5000, p: 14)
:ok = ExDataSketch.GenStage.SketchConsumer.merge(consumer, partial)

# Check the estimate
estimate = ExDataSketch.GenStage.SketchConsumer.estimate(consumer)
IO.puts("Consumer estimate: #{Float.round(estimate, 0)} unique users (true: 5000)")

# Get the full sketch
sketch = ExDataSketch.GenStage.SketchConsumer.get(consumer)
IO.puts("Full sketch estimate: #{Float.round(ExDataSketch.HLL.estimate(sketch), 0)} (true: 5000)")

GenServer.stop(consumer)

Section 2: Flushing and Resetting

Flush returns the accumulated sketch and resets the consumer with a fresh empty sketch. This is useful for rolling window aggregation.

{:ok, consumer} = ExDataSketch.GenStage.SketchConsumer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14],
  subscribe_to: []
)

# Merge some data
partial1 = ExDataSketch.HLL.from_enumerable(1..2000, p: 14)
:ok = ExDataSketch.GenStage.SketchConsumer.merge(consumer, partial1)

# Flush: get sketch and reset
flushed = ExDataSketch.GenStage.SketchConsumer.flush(consumer)
IO.puts("Flushed estimate: #{Float.round(ExDataSketch.HLL.estimate(flushed), 0)} (true: 2000)")

# After flush, the consumer starts fresh
estimate = ExDataSketch.GenStage.SketchConsumer.estimate(consumer)
IO.puts("Post-flush estimate: #{estimate} (should be 0.0)")

GenServer.stop(consumer)

Section 3: SketchProducer -- Emit Sketches on Demand

SketchProducer is a GenStage producer that emits its accumulated sketch when consumers demand events. Downstream consumers receive sketch snapshots.

# Start a producer that accumulates items
{:ok, producer} = ExDataSketch.GenStage.SketchProducer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14]
)

# Update the producer with items
:ok = ExDataSketch.GenStage.SketchProducer.update(producer, "user_1")
:ok = ExDataSketch.GenStage.SketchProducer.update(producer, "user_2")
:ok = ExDataSketch.GenStage.SketchProducer.update(producer, "user_3")

# Merge a partial sketch
partial = ExDataSketch.HLL.from_enumerable(4..100, p: 14)
:ok = ExDataSketch.GenStage.SketchProducer.merge(producer, partial)

IO.puts("Producer estimate: #{Float.round(ExDataSketch.GenStage.SketchProducer.estimate(producer), 0)} (true: 100)")

GenServer.stop(producer)

Section 4: Periodic Flushing with Callback

SketchConsumer supports automatic periodic flushing with a callback, ideal for pushing sketches to telemetry or persistence on a timer:

# Consumer that flushes every 5 seconds and sends estimate to telemetry
{:ok, consumer} = ExDataSketch.GenStage.SketchConsumer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14],
  flush_interval: 5_000,
  flush_callback: fn sketch ->
    estimate = ExDataSketch.HLL.estimate(sketch)
    IO.puts("[FLUSH CALLBACK] Estimate: #{Float.round(estimate, 0)} (true: 1000)")
  end,
  subscribe_to: []
)

# Merge data -- in a real pipeline this comes from a producer
partial = ExDataSketch.HLL.from_enumerable(1..1000, p: 14)
:ok = ExDataSketch.GenStage.SketchConsumer.merge(consumer, partial)

# The flush_callback will fire every 5 seconds
# For demo, trigger a manual flush
flushed = ExDataSketch.GenStage.SketchConsumer.flush(consumer)
IO.puts("Manual flush estimate: #{Float.round(ExDataSketch.HLL.estimate(flushed), 0)} (true: 1000)")

GenServer.stop(consumer)

Section 5: Putting It Together -- Producer to Consumer Pipeline

SketchProducer emits sketch snapshots of its current accumulated sketch. When SketchConsumer subscribes to a SketchProducer, each event it receives is a sketch struct of the configured :sketch_module and is merged into the consumer's accumulated sketch via sketch_module.merge/2.

Raw event sources (Kafka, RabbitMQ, etc.) emit non-sketch terms; the consumer detects them by struct type and instead applies :key_fn and update/2. A single batch can mix both shapes.

# Producer that accumulates items and emits snapshots downstream.
{:ok, producer} = ExDataSketch.GenStage.SketchProducer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14]
)

# Consumer that merges incoming snapshots into its own sketch.
{:ok, consumer} = ExDataSketch.GenStage.SketchConsumer.start_link(
  sketch_module: ExDataSketch.HLL,
  sketch_opts: [p: 14],
  subscribe_to: [{producer, max_demand: 100}]
)

# Feed the producer; each update triggers one snapshot to the consumer.
items = Enum.map(1..500, &"user_#{&1}")
Enum.each(items, fn item -> :ok = ExDataSketch.GenStage.SketchProducer.update(producer, item) end)

# Allow the demand/emit loop to settle.
Process.sleep(200)

producer_estimate = ExDataSketch.GenStage.SketchProducer.estimate(producer)
consumer_estimate = ExDataSketch.GenStage.SketchConsumer.estimate(consumer)

IO.puts("Producer estimate: #{Float.round(producer_estimate, 0)} (true: 500)")
IO.puts("Consumer estimate: #{Float.round(consumer_estimate, 0)} (true: 500)")
IO.puts("Round-trip preserved cardinality: #{abs(consumer_estimate - 500) < 25}")

GenServer.stop(consumer)
GenServer.stop(producer)

Section 6: Operational Guidance

Use GenStage when:

  • You need back-pressure from consumers to producers
  • Data arrives continuously from an external source (Kafka, RabbitMQ)
  • You want bounded memory usage (sketch is O(2^p) regardless of input size)
  • You need periodic snapshots of evolving sketches

Use Broadway instead when:

  • You need batch processing, acknowledgements, or rate limiting
  • You want built-in concurrency and partition handling
  • You need dead-letter handling and retry logic

Memory: At p=14, both producer and consumer hold a 16KB sketch. 4 partitions = 64KB total. Flushing resets the sketch each window.

Flush interval: Choose based on latency requirements:

  • 1-5 seconds for real-time dashboards
  • 30-60 seconds for batch analytics
  • :infinity (default) for manual flush only