code wiki / (root) / nx_distexec.nx

nx_distexec.nx source

↩ module page · 208 lines · 8513 B

1// nx_distexec.nx -- DISTRIBUTED EXECUTION ALGEBRA (distributed-compute, F-dist). Spark's essence is not 2// the cluster -- it is the PARTITION -> MAP -> MERGE algebra: split the data into partitions, compute a 3// PARTIAL result per partition, then MERGE the partials into the whole answer. That algebra is provable and 4// buildable single-process from the first byte up; only its PARALLEL execution (a partition per worker 5// core/node) needs the threading/swarm arc. This organ is the algebra, proven PARTITION-INVARIANT: the 6// merged result is byte-identical for ANY partition count P, which is exactly the correctness a scheduler 7// relies on to run partitions anywhere. Composes nx_vecexec (each partition aggregates vectorized) over the 8// zero-copy nx_colframe. Honest: this is the algebra + a sequential executor; real multi-core/multi-node 9// parallelism is the named finer gap. license_tier: ORIGINAL No hardware writes (Rule 26). 10import "nx_syscalls.nx" 11import "nx_colframe.nx" 12const DE_MAGIC_4242: i64 = 4242 13const DE_MAGIC_1103515245: i64 = 1103515245 14const DE_MAGIC_12345: i64 = 12345 15const DE_MAGIC_100000: i64 = 100000 16const DE_MAGIC_4000000: i64 = 4000000 17 18const DE_VEC: i64 = 1024 19const DE_MAXK: i64 = 256 // group-by key domain cap (buckets) 20const DE_OUT: i64 = 8192 21const DE_ZERO: i64 = 48 22// a partial aggregate is [sum, count, min, max] 23const DE_SUM: i64 = 0 24const DE_CNT: i64 = 1 25const DE_MIN: i64 = 2 26const DE_MAX: i64 = 3 27 28// ---- MAP: partial aggregate over ONE partition (row range [lo,hi)), vectorized within the partition 29func de_agg_range(frame: *u8, ci: i64, lo: i64, hi: i64, out4: *i64) -> i64 { 30 let col: *i64 = cf_col(frame, ci) 31 var sum: i64 = 0 32 var cnt: i64 = 0 33 var mn: i64 = 0 34 var mx: i64 = 0 35 var first: i64 = 1 36 var base: i64 = lo 37 while base < hi { 38 var lim: i64 = base + DE_VEC 39 if lim > hi { lim = hi } 40 var r: i64 = base 41 while r < lim { 42 let v: i64 = col[r] 43 sum = sum + v 44 cnt = cnt + 1 45 if first == 1 { mn = v; mx = v; first = 0 } else { if v < mn { mn = v } if v > mx { mx = v } } 46 r = r + 1 47 } 48 base = lim 49 } 50 out4[DE_SUM] = sum 51 out4[DE_CNT] = cnt 52 out4[DE_MIN] = mn 53 out4[DE_MAX] = mx 54 return 0 55} 56// ---- MERGE: combine partial b INTO accumulator a (associative + commutative => any partition order) ---- 57func de_merge(a: *i64, b: *i64) -> i64 { 58 if b[DE_CNT] == 0 { return 0 } 59 if a[DE_CNT] == 0 { 60 a[DE_SUM] = b[DE_SUM] 61 a[DE_CNT] = b[DE_CNT] 62 a[DE_MIN] = b[DE_MIN] 63 a[DE_MAX] = b[DE_MAX] 64 return 0 65 } 66 a[DE_SUM] = a[DE_SUM] + b[DE_SUM] 67 a[DE_CNT] = a[DE_CNT] + b[DE_CNT] 68 if b[DE_MIN] < a[DE_MIN] { a[DE_MIN] = b[DE_MIN] } 69 if b[DE_MAX] > a[DE_MAX] { a[DE_MAX] = b[DE_MAX] } 70 return 0 71} 72// partition boundary: the start row of partition p of P over n rows (balanced). 73func de_bound(n: i64, P: i64, p: i64) -> i64 { 74 // floor(n*p/P) gives balanced, gapless, covering boundaries 75 return n * p / P 76} 77// ---- the DISTRIBUTED JOB: aggregate `ci` over P partitions, merge to one result in out4 -------------- 78func de_run_agg(frame: *u8, ci: i64, P: i64, out4: *i64) -> i64 { 79 let n: i64 = cf_nrows(frame) 80 var pp: i64 = P 81 if pp < 1 { pp = 1 } 82 out4[DE_SUM] = 0 83 out4[DE_CNT] = 0 84 out4[DE_MIN] = 0 85 out4[DE_MAX] = 0 86 let part: *i64 = sys_mmap(8 * 4) as *i64 87 var p: i64 = 0 88 while p < pp { 89 let lo: i64 = de_bound(n, pp, p) 90 let hi: i64 = de_bound(n, pp, p + 1) 91 de_agg_range(frame, ci, lo, hi, part) 92 de_merge(out4, part) 93 p = p + 1 94 } 95 return 0 96} 97 98// ---- GROUP BY key SUM, the shuffle-merge case: per-partition bucket sums, merged bucket-wise ---------- 99func de_gb_range(frame: *u8, keyci: i64, valci: i64, lo: i64, hi: i64, buckets: *i64, K: i64) -> i64 { 100 let kc: *i64 = cf_col(frame, keyci) 101 let vc: *i64 = cf_col(frame, valci) 102 var r: i64 = lo 103 while r < hi { 104 let k: i64 = kc[r] 105 if k >= 0 { if k < K { buckets[k] = buckets[k] + vc[r] } } 106 r = r + 1 107 } 108 return 0 109} 110func de_run_groupby(frame: *u8, keyci: i64, valci: i64, K: i64, P: i64, buckets: *i64) -> i64 { 111 let n: i64 = cf_nrows(frame) 112 var pp: i64 = P 113 if pp < 1 { pp = 1 } 114 var i: i64 = 0 115 while i < K { buckets[i] = 0; i = i + 1 } 116 // each partition sums into its OWN bucket set (the map side), then we merge bucket-wise (the shuffle) 117 let pb: *i64 = sys_mmap(8 * K) as *i64 118 var p: i64 = 0 119 while p < pp { 120 var j: i64 = 0 121 while j < K { pb[j] = 0; j = j + 1 } 122 let lo: i64 = de_bound(n, pp, p) 123 let hi: i64 = de_bound(n, pp, p + 1) 124 de_gb_range(frame, keyci, valci, lo, hi, pb, K) 125 var m: i64 = 0 126 while m < K { buckets[m] = buckets[m] + pb[m]; m = m + 1 } 127 p = p + 1 128 } 129 return 0 130} 131 132// ---- a small live demo/report over a stored dataset (agentic JSON) -------------------------------- 133func de_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 } 134func de_num(out: *u8, o: i64, v: i64) -> i64 { 135 var oo: i64 = o 136 var m: i64 = v 137 if m < 0 { out[oo] = 45 as u8; oo = oo + 1; m = 0 - m } 138 let t: *u8 = sys_mmap(24) 139 var k: i64 = 0 140 if m == 0 { t[0] = DE_ZERO as u8; k = 1 } 141 while m > 0 { t[k] = (DE_ZERO + (m % 10)) as u8; m = m / 10; k = k + 1 } 142 var i: i64 = k - 1 143 while i >= 0 { out[oo] = t[i]; oo = oo + 1; i = i - 1 } 144 return oo 145} 146func de_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 } 147 148// build a demo frame and show the SAME aggregate under P=1 and P=P, proving partition-invariance live 149func de_demo(nrows: i64, P: i64, out: *u8) -> i64 { 150 let val: *i64 = sys_mmap(8 * nrows) as *i64 151 var s: i64 = DE_MAGIC_4242 152 var i: i64 = 0 153 while i < nrows { s = (s * DE_MAGIC_1103515245 + DE_MAGIC_12345) & 0x7fffffff; val[i] = s % DE_MAGIC_100000; i = i + 1 } 154 let cols: *i64 = sys_mmap(8 * 1) as *i64 155 let names: *i64 = sys_mmap(8 * 1) as *i64 156 cols[0] = val as i64 157 names[0] = "val" as *u8 as i64 158 let tot: i64 = cf_encoded_bytes(names, 1, nrows) 159 let frame: *u8 = sys_mmap(tot + 64) 160 cf_encode(cols, names, 1, nrows, frame) 161 let a1: *i64 = sys_mmap(8 * 4) as *i64 162 let aP: *i64 = sys_mmap(8 * 4) as *i64 163 de_run_agg(frame, 0, 1, a1) 164 de_run_agg(frame, 0, P, aP) 165 var o: i64 = 0 166 o = de_raw(out, o, "{\"tool\":\"nx_distexec\",\"verb\":\"demo\",\"rows\":" as *u8) 167 o = de_num(out, o, nrows) 168 o = de_raw(out, o, ",\"partitions\":" as *u8) 169 o = de_num(out, o, P) 170 o = de_raw(out, o, ",\"agg\":\"sum/count/min/max(val)\",\"single_partition\":{\"sum\":" as *u8) 171 o = de_num(out, o, a1[DE_SUM]) 172 o = de_raw(out, o, ",\"count\":" as *u8) 173 o = de_num(out, o, a1[DE_CNT]) 174 o = de_raw(out, o, ",\"min\":" as *u8) 175 o = de_num(out, o, a1[DE_MIN]) 176 o = de_raw(out, o, ",\"max\":" as *u8) 177 o = de_num(out, o, a1[DE_MAX]) 178 o = de_raw(out, o, "},\"partitioned\":{\"sum\":" as *u8) 179 o = de_num(out, o, aP[DE_SUM]) 180 o = de_raw(out, o, ",\"count\":" as *u8) 181 o = de_num(out, o, aP[DE_CNT]) 182 o = de_raw(out, o, ",\"min\":" as *u8) 183 o = de_num(out, o, aP[DE_MIN]) 184 o = de_raw(out, o, ",\"max\":" as *u8) 185 o = de_num(out, o, aP[DE_MAX]) 186 o = de_raw(out, o, "},\"partition_invariant\":" as *u8) 187 var inv: i64 = 1 188 if a1[DE_SUM] != aP[DE_SUM] { inv = 0 } 189 if a1[DE_CNT] != aP[DE_CNT] { inv = 0 } 190 if a1[DE_MIN] != aP[DE_MIN] { inv = 0 } 191 if a1[DE_MAX] != aP[DE_MAX] { inv = 0 } 192 if inv == 1 { o = de_raw(out, o, "true" as *u8) } else { o = de_raw(out, o, "false" as *u8) } 193 o = de_raw(out, o, ",\"note\":\"merged result is identical whichever partition count runs it -- a scheduler may place partitions on any worker; parallel execution is the finer gap\"}\n" as *u8) 194 out[o] = 0 as u8 195 return o 196} 197func main(argc: i64, argv: *i64) -> i64 { 198 let out: *u8 = sys_mmap(DE_OUT) 199 var nrows: i64 = DE_MAGIC_4000000 200 var P: i64 = 8 201 if argc >= 2 { nrows = de_atoi(argv[1] as *u8) } 202 if argc >= 3 { P = de_atoi(argv[2] as *u8) } 203 if nrows < 1 { nrows = DE_MAGIC_4000000 } 204 if P < 1 { P = 8 } 205 let n: i64 = de_demo(nrows, P, out) 206 sys_write(1, out, n) 207 return 0 208}