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}