Level 1 Global mode across shards
A dataset of integers is split across the workers of a cluster. Links are slow, so what counts is how much you send, not local CPU. You coordinate the workers through a cluster object:
cluster.num_workers: number of workersW.cluster.map(fn, *args) -> list: runsfn(worker, *args)on every worker and returns the results to you, in worker-id order.cluster.shuffle(fn, *args) -> None: runsfn(worker, *args)on every worker. It must return a dict{target_worker_id: [items...]}. Every worker'sinboxis emptied, then each item is delivered to its target'sinbox.- A
workerhasworker.id,worker.num_workers,worker.data(its shard: a list of ints; don't modify it) andworker.inbox(a list).
Traffic is measured in numbers: each int, float or string inside a payload counts 1, however deeply it's nested in lists, tuples or dicts (dict keys count too). None is free.
- Every result a
mapreturns counts as coordinator traffic, and so do theargsyou pass, once per worker (it's a broadcast). - Shuffle items sent to another worker count as peer traffic. Items a worker sends to itself are free.
fnmust be a plain function with no closure over your local variables (the cluster raisesValueErrorotherwise). Pass what it needs throughargs, where it is paid for.
Implement global_mode(cluster) -> int: the value with the highest total count across all shards, with ties going to the smallest value. At least one shard is non-empty.
Budget: total traffic at most 2 * D + 4 * W, where D is the sum over workers of the number of distinct values in that worker's shard. Shipping every raw value to yourself blows this budget when shards contain repeats. Sending only each worker's top few values is cheap but wrong: the global mode need not be near the top on any single worker.
# 2 workers: [5, 5, 1, 1, 1] and [5, 5, 2, 2, 2]
global_mode(cluster) # 5 (count 4), even though it wins on neither worker