code wiki / (root) / nx_mrshuffle.nx

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}