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.