nx_mrshuffle.nx source
↩ module page · 170 lines · 7517 B
1// nx_mrshuffle.nx -- the real MapReduce SHUFFLE (closes distributed-compute's shuffle-exchange gap). The
2// hard, distinctive part of distributed group-by/join is the SHUFFLE: repartition rows by a hash of the
3// key so that every row for a given key lands at exactly ONE reducer, and each reducer sees ALL rows for
4// its keys -- then reducers aggregate independently and their union is the answer. This is the canonical
5// COUNTING shuffle (count per reducer -> prefix-sum to contiguous offsets -> place), so it is O(n) not the
6// O(n*R) broadcast-filter shortcut, and the reduce phase runs in REAL parallel (one forked process per
7// reducer, over a shared exchange). From the first byte up: fork + shared mmap + wait4, no frameworks.
8// Composes the zero-copy nx_colframe. license_tier: ORIGINAL No hardware writes (Rule 26).
9import "nx_syscalls.nx"
10import "nx_colframe.nx"
11const MR_MAGIC_4294967296: i64 = 4294967296
12const MR_MAGIC_1103515245: i64 = 1103515245
13const MR_MAGIC_12345: i64 = 12345
14const MR_MAGIC_1000000: i64 = 1000000
15
16const MR_MAXR: i64 = 64
17const MR_OUT: i64 = 8192
18const MR_ZERO: i64 = 48
19const MR_HASHK: i64 = 2654435761 // Knuth multiplicative hash constant (odd)
20const MR_MASK: i64 = 0x7fffffffffffffff
21
22func mr_hash(k: i64) -> i64 {
23 var h: i64 = k
24 if h < 0 { h = 0 - h }
25 h = (h * MR_HASHK) & MR_MASK
26 h = (h ^ (h / MR_MAGIC_4294967296)) & MR_MASK // mix high bits down (>>32 via divide, no shift op needed)
27 return h & MR_MASK
28}
29func mr_reducer(k: i64, R: i64) -> i64 { return mr_hash(k) % R }
30
31// GROUP BY key SUM via a real shuffle + parallel reduce. keys are assumed in [0,K) for a key-indexed
32// output (out_sums must be a SHARED array of K i64); arbitrary keys = a per-reducer hash table, the named
33// refinement. out_sums is written by the reducers (each owns a disjoint key set -> disjoint writes, no
34// locks). Returns 1 on success, 0 on capacity refusal.
35func mr_groupby_shuffle(frame: *u8, keyci: i64, valci: i64, R: i64, out_sums: *i64) -> i64 {
36 let n: i64 = cf_nrows(frame)
37 var rr: i64 = R
38 if rr < 1 { rr = 1 }
39 if rr > MR_MAXR { rr = MR_MAXR }
40 let kc: *i64 = cf_col(frame, keyci)
41 let vc: *i64 = cf_col(frame, valci)
42
43 // --- SHUFFLE step 1: count rows per reducer (map-side histogram) ---
44 let cnt: *i64 = sys_mmap(8 * rr) as *i64
45 var i: i64 = 0
46 while i < rr { cnt[i] = 0; i = i + 1 }
47 var r: i64 = 0
48 while r < n { let red: i64 = mr_reducer(kc[r], rr); cnt[red] = cnt[red] + 1; r = r + 1 }
49 // --- step 2: prefix-sum to contiguous per-reducer offsets ---
50 let off: *i64 = sys_mmap(8 * (rr + 1)) as *i64
51 off[0] = 0
52 i = 0
53 while i < rr { off[i + 1] = off[i] + cnt[i]; i = i + 1 }
54 // --- step 3: place each (key,val) into the SHARED exchange at its reducer's cursor ---
55 let ex: *i64 = sys_mmap_shared(8 * 2 * (n + 1)) as *i64
56 let cur: *i64 = sys_mmap(8 * rr) as *i64
57 i = 0
58 while i < rr { cur[i] = off[i]; i = i + 1 }
59 r = 0
60 while r < n {
61 let red: i64 = mr_reducer(kc[r], rr)
62 let p: i64 = cur[red]
63 ex[p * 2] = kc[r]
64 ex[p * 2 + 1] = vc[r]
65 cur[red] = p + 1
66 r = r + 1
67 }
68 // --- step 4: PARALLEL REDUCE -- fork one process per reducer over its contiguous exchange slice ---
69 let pids: *i64 = sys_mmap(8 * rr) as *i64
70 var nwait: i64 = 0
71 var q: i64 = 0
72 while q < rr {
73 let lo: i64 = off[q]
74 let hi: i64 = off[q + 1]
75 let pid: i64 = sys_fork()
76 if pid == 0 {
77 var p: i64 = lo
78 while p < hi { let k: i64 = ex[p * 2]; out_sums[k] = out_sums[k] + ex[p * 2 + 1]; p = p + 1 }
79 sys_exit(0)
80 }
81 if pid > 0 { pids[nwait] = pid; nwait = nwait + 1 } else {
82 var p2: i64 = lo
83 while p2 < hi { let k: i64 = ex[p2 * 2]; out_sums[k] = out_sums[k] + ex[p2 * 2 + 1]; p2 = p2 + 1 }
84 }
85 q = q + 1
86 }
87 let st: *i64 = sys_mmap(16) as *i64
88 var w: i64 = 0
89 while w < nwait { sys_wait4(pids[w], st, 0); w = w + 1 }
90 return 1
91}
92
93// ---- live demo: shuffle group-by over a generated dataset, report the exchange + a few group sums -----
94func mr_raw(out: *u8, o: i64, s: *u8) -> i64 { var i: i64 = 0; while s[i] != (0 as u8) { out[o] = s[i]; o = o + 1; i = i + 1 } return o }
95func mr_num(out: *u8, o: i64, v: i64) -> i64 {
96 var oo: i64 = o
97 var m: i64 = v
98 if m < 0 { out[oo] = 45 as u8; oo = oo + 1; m = 0 - m }
99 let t: *u8 = sys_mmap(24)
100 var k: i64 = 0
101 if m == 0 { t[0] = MR_ZERO as u8; k = 1 }
102 while m > 0 { t[k] = (MR_ZERO + (m % 10)) as u8; m = m / 10; k = k + 1 }
103 var i: i64 = k - 1
104 while i >= 0 { out[oo] = t[i]; oo = oo + 1; i = i - 1 }
105 return oo
106}
107func mr_atoi(s: *u8) -> i64 { var v: i64 = 0; var i: i64 = 0; var go: i64 = 1; while go == 1 { let c: i64 = s[i] as i64; if c < 48 { go = 0 } else { if c > 57 { go = 0 } else { v = v * 10 + (c - 48); i = i + 1 } } } return v }
108
109func mr_demo(nrows: i64, K: i64, R: i64, out: *u8) -> i64 {
110 let key: *i64 = sys_mmap(8 * nrows) as *i64
111 let val: *i64 = sys_mmap(8 * nrows) as *i64
112 var s: i64 = 99
113 var i: i64 = 0
114 while i < nrows { s = (s * MR_MAGIC_1103515245 + MR_MAGIC_12345) & 0x7fffffff; key[i] = s % K; val[i] = (s / K) % 1000; i = i + 1 }
115 let cols: *i64 = sys_mmap(8 * 2) as *i64
116 let names: *i64 = sys_mmap(8 * 2) as *i64
117 cols[0] = key as i64
118 cols[1] = val as i64
119 names[0] = "key" as *u8 as i64
120 names[1] = "val" as *u8 as i64
121 let tot: i64 = cf_encoded_bytes(names, 2, nrows)
122 let frame: *u8 = sys_mmap(tot + 64)
123 cf_encode(cols, names, 2, nrows, frame)
124
125 let shuf: *i64 = sys_mmap_shared(8 * K) as *i64
126 var z: i64 = 0
127 while z < K { shuf[z] = 0; z = z + 1 }
128 mr_groupby_shuffle(frame, 0, 1, R, shuf)
129 // reference: monolithic group-by
130 let ref: *i64 = sys_mmap(8 * K) as *i64
131 z = 0
132 while z < K { ref[z] = 0; z = z + 1 }
133 i = 0
134 while i < nrows { ref[key[i]] = ref[key[i]] + val[i]; i = i + 1 }
135 var agree: i64 = 1
136 z = 0
137 while z < K { if shuf[z] != ref[z] { agree = 0 } z = z + 1 }
138
139 var o: i64 = 0
140 o = mr_raw(out, o, "{\"tool\":\"nx_mrshuffle\",\"verb\":\"demo\",\"rows\":" as *u8)
141 o = mr_num(out, o, nrows)
142 o = mr_raw(out, o, ",\"distinct_keys\":" as *u8)
143 o = mr_num(out, o, K)
144 o = mr_raw(out, o, ",\"reducers\":" as *u8)
145 o = mr_num(out, o, R)
146 o = mr_raw(out, o, ",\"shuffle_groupby_matches_monolithic\":" as *u8)
147 if agree == 1 { o = mr_raw(out, o, "true" as *u8) } else { o = mr_raw(out, o, "false" as *u8) }
148 o = mr_raw(out, o, ",\"sample_group_sums\":{\"key0\":" as *u8)
149 o = mr_num(out, o, shuf[0])
150 o = mr_raw(out, o, ",\"key1\":" as *u8)
151 o = mr_num(out, o, shuf[1])
152 o = mr_raw(out, o, "},\"note\":\"counting shuffle (count->prefix-sum->place, O(n)) + parallel per-reducer reduce over a shared exchange; fork+shared-mmap+wait4, no frameworks. arbitrary (unbounded) keys via per-reducer hash tables is the finer step.\"}\n" as *u8)
153 out[o] = 0 as u8
154 return o
155}
156func main(argc: i64, argv: *i64) -> i64 {
157 let out: *u8 = sys_mmap(MR_OUT)
158 var nrows: i64 = MR_MAGIC_1000000
159 var K: i64 = 100
160 var R: i64 = 4
161 if argc >= 2 { nrows = mr_atoi(argv[1] as *u8) }
162 if argc >= 3 { K = mr_atoi(argv[2] as *u8) }
163 if argc >= 4 { R = mr_atoi(argv[3] as *u8) }
164 if nrows < 1 { nrows = MR_MAGIC_1000000 }
165 if K < 1 { K = 100 }
166 if R < 1 { R = 4 }
167 let n: i64 = mr_demo(nrows, K, R, out)
168 sys_write(1, out, n)
169 return 0
170}