Mix.install([
  {:ex_data_sketch, "~> 0.10"}
])

Introduction

Sketch merging is associative and commutative, making it ideal for distributed aggregation. This Livebook demonstrates patterns for merging sketches across processes, nodes, and time windows.

Section 1: Associativity and Commutativity

# Generate three independent data sets
items_a = for i <- 1..3000, do: "user_#{i}"
items_b = for i <- 2000..5000, do: "user_#{i}"
items_c = for i <- 4000..7000, do: "user_#{i}"

# Create sketches from each
sketch_a = ExDataSketch.HLL.from_enumerable(items_a, p: 14)
sketch_b = ExDataSketch.HLL.from_enumerable(items_b, p: 14)
sketch_c = ExDataSketch.HLL.from_enumerable(items_c, p: 14)

true_count = MapSet.new(items_a ++ items_b ++ items_c) |> MapSet.size()
IO.puts("True distinct count: #{true_count}")

# Associativity: (a + b) + c == a + (b + c)
left = ExDataSketch.HLL.merge(ExDataSketch.HLL.merge(sketch_a, sketch_b), sketch_c)
right = ExDataSketch.HLL.merge(sketch_a, ExDataSketch.HLL.merge(sketch_b, sketch_c))

IO.puts("Left associative: #{Float.round(ExDataSketch.HLL.estimate(left), 0)} (true: #{true_count})")
IO.puts("Right associative: #{Float.round(ExDataSketch.HLL.estimate(right), 0)} (true: #{true_count})")

# Commutativity: a + b == b + a
ab = ExDataSketch.HLL.merge(sketch_a, sketch_b)
ba = ExDataSketch.HLL.merge(sketch_b, sketch_a)
IO.puts("a+b = #{Float.round(ExDataSketch.HLL.estimate(ab), 0)}, b+a = #{Float.round(ExDataSketch.HLL.estimate(ba), 0)} (true: ~5000)")

# merge_many merges all at once
all = ExDataSketch.HLL.merge_many([sketch_a, sketch_b, sketch_c])
IO.puts("merge_many: #{Float.round(ExDataSketch.HLL.estimate(all), 0)} (true: #{true_count})")

Section 2: Multi-Node Aggregation Pattern

# Simulate 5 independent nodes, each processing a partition
partitions = Enum.chunk_every(1..50_000, 10_000)

node_sketches = Enum.map(partitions, fn partition ->
  items = Enum.map(partition, fn i -> "user_#{i}" end)
  ExDataSketch.HLL.from_enumerable(items, p: 14)
end)

# Each node sends its sketch to the aggregator
merged = ExDataSketch.HLL.merge_many(node_sketches)
IO.puts("5-node merge estimate: #{Float.round(ExDataSketch.HLL.estimate(merged), 0)} (true: 50000)")
IO.puts("Memory sent: #{5 * ExDataSketch.HLL.size_bytes(hd(node_sketches))} bytes (5 x #{ExDataSketch.HLL.size_bytes(hd(node_sketches))} bytes)")

Section 3: Hierarchical Aggregation

For very large clusters, use tree aggregation to reduce the merge burden:

# 16 leaf nodes, 4 intermediate aggregators, 1 root
leaves = for group <- 0..3 do
  partition = Enum.map(1..2500, fn i -> "node_#{group}_user_#{i}" end)
  ExDataSketch.HLL.from_enumerable(partition, p: 14)
end

# Intermediate level: merge 4 leaves each
intermediates = leaves
  |> Enum.chunk_every(4)
  |> Enum.map(&ExDataSketch.HLL.merge_many/1)

# Root: merge 4 intermediates
root = ExDataSketch.HLL.merge_many(intermediates)
IO.puts("Tree-aggregated estimate: #{Float.round(ExDataSketch.HLL.estimate(root), 0)} (true: 10000)")

# Compare with flat merge
flat = ExDataSketch.HLL.merge_many(leaves)
IO.puts("Flat merge estimate: #{Float.round(ExDataSketch.HLL.estimate(flat), 0)} (true: 10000)")

Section 4: ETS-Sharded Aggregation

Use ETS for shared sketch state accessible from any process:

table = :ets.new(:distributed_sketches, [:set, :public, :named_table])

# Simulate 3 workers each saving to ETS
worker_a = ExDataSketch.HLL.from_enumerable(Enum.map(1..3000, fn i -> "user_#{i}" end), p: 14)
worker_b = ExDataSketch.HLL.from_enumerable(Enum.map(2000..5000, fn i -> "user_#{i}" end), p: 14)

ExDataSketch.Storage.ETS.save(worker_a, table, "region:a")
ExDataSketch.Storage.ETS.merge(worker_b, table, "region:a")

{:ok, merged} = ExDataSketch.Storage.ETS.load(ExDataSketch.HLL, table, "region:a")
IO.puts("ETS merged estimate: #{Float.round(ExDataSketch.HLL.estimate(merged), 0)} (true: 5000)")

:ets.delete(table)

Section 5: Operational Guidance

Merge frequency: Merge every N seconds (N = 1-60 depending on latency requirements). More frequent = lower latency but more network traffic.

Sketch size at scale: At p=14, each sketch = 16KB. Merging 1000 sketches per second = 16MB/s network bandwidth. At p=10, it's only 1MB/s.

Consistency: Sketch merging provides eventual consistency. During network partitions, partitions accumulate independently. When connectivity is restored, simply merge all partition sketches.