Files
wickra/examples/node/parallel_assets.js
T
kingchenc 43b0b26736 examples: add parallel-assets demos for Rust and Node
Python's parallel_assets.py demoed GIL-release multi-core throughput;
Rust and Node both lacked a sibling that shows their own native
parallelism. Close the gap with two real, runnable examples.

* examples/rust/src/bin/parallel_assets.rs — synthesises an (assets,
  bars) panel with a deterministic per-asset LCG, runs a serial baseline,
  then `Sma::batch_parallel` / `Rsi::batch_parallel` via rayon, asserts
  the two outputs are element-wise identical and prints the speedup.
  Toggle indicator with `--indicator sma|rsi`.
* examples/node/parallel_assets.js — same shape, but the parallel run is
  a `worker_threads` pool that re-loads the native binding in each
  worker. Each worker computes the last non-null indicator value for its
  slice; the main thread aggregates and verifies serial == parallel
  per asset.

Both examples report timings and the serial-vs-parallel sanity check
passes. Defaults (200 × 5000) keep the example fast on dev hardware;
larger `--assets`/`--bars` is where the speedup numbers move (Node's
worker spawn cost dominates the smallest sizes, which is honest and
educational).

examples/README.md gains the two new rows.
2026-05-23 00:18:46 +02:00

170 lines
5.8 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Parallel multi-asset indicator computation via Node.js worker_threads.
//
// Builds a synthetic (assets, bars) panel, runs a serial baseline on the
// main thread, then dispatches the same workload to a pool of workers
// (each loading its own copy of the native wickra binding) and reports
// the speedup. The Node counterpart of examples/python/parallel_assets.py.
//
// Run with:
// node parallel_assets.js [--assets N] [--bars M]
// [--indicator sma|rsi] [--workers W]
const { Worker, isMainThread, workerData, parentPort } = require('node:worker_threads');
const os = require('node:os');
// Deterministic LCG matching the Rust sibling so a side-by-side serial
// computation on both languages produces visibly comparable timings.
function makeRng(seed) {
let state = seed >>> 0;
return () => {
state = (Math.imul(state, 1103515245) + 12345) & 0x7fffffff;
return state / 0x7fffffff;
};
}
function synthesizePanel(nAssets, nBars) {
const panel = new Array(nAssets);
for (let a = 0; a < nAssets; a++) {
const rng = makeRng((1234567 + Math.imul(a, 2654435761)) >>> 0);
const series = new Float64Array(nBars);
let price = 100.0;
for (let i = 0; i < nBars; i++) {
price += (rng() - 0.5) * 0.4;
series[i] = price;
}
panel[a] = series;
}
return panel;
}
// Compute the last non-null indicator output for one price series. Kept
// tiny so the same code runs on the main thread (serial baseline) and
// inside each worker (parallel run).
function lastValue(prices, indicator, wickra) {
const ind = indicator === 'sma' ? new wickra.SMA(14) : new wickra.RSI(14);
let last = null;
for (let i = 0; i < prices.length; i++) {
const v = ind.update(prices[i]);
if (v !== null) last = v;
}
return last;
}
if (!isMainThread) {
// ---------------- worker side ----------------
const wickra = require('wickra');
const { panelSlice, indicator } = workerData;
const results = new Array(panelSlice.length);
for (let i = 0; i < panelSlice.length; i++) {
results[i] = lastValue(panelSlice[i], indicator, wickra);
}
parentPort.postMessage(results);
} else {
// ---------------- main side ----------------
const wickra = require('wickra');
function parseArgs(argv) {
const args = {
assets: 200,
bars: 5000,
indicator: 'sma',
workers: Math.max(1, os.cpus().length),
};
for (let i = 0; i < argv.length; i++) {
const k = argv[i];
if (k === '--assets') args.assets = Number(argv[++i]);
else if (k === '--bars') args.bars = Number(argv[++i]);
else if (k === '--indicator') args.indicator = argv[++i];
else if (k === '--workers') args.workers = Number(argv[++i]);
else throw new Error(`unexpected argument: ${k}`);
}
if (!Number.isInteger(args.assets) || args.assets <= 0) {
throw new Error('--assets must be a positive integer');
}
if (!Number.isInteger(args.bars) || args.bars <= 0) {
throw new Error('--bars must be a positive integer');
}
if (args.indicator !== 'sma' && args.indicator !== 'rsi') {
throw new Error("--indicator: expected 'sma' or 'rsi'");
}
if (!Number.isInteger(args.workers) || args.workers <= 0) {
throw new Error('--workers must be a positive integer');
}
return args;
}
let args;
try {
args = parseArgs(process.argv.slice(2));
} catch (err) {
console.error(`error: ${err.message}`);
process.exit(2);
}
console.log(`Generating ${args.assets}×${args.bars} synthetic panel…`);
const panel = synthesizePanel(args.assets, args.bars);
// Serial baseline (main thread).
let t0 = process.hrtime.bigint();
const serial = new Array(args.assets);
for (let a = 0; a < args.assets; a++) {
serial[a] = lastValue(panel[a], args.indicator, wickra);
}
const tSerial = Number(process.hrtime.bigint() - t0) / 1e9;
console.log(
`Serial: ${tSerial.toFixed(3).padStart(8)} s (${args.assets} assets, indicator=${args.indicator})`,
);
// Parallel via worker_threads.
const workerCount = Math.max(1, Math.min(args.workers, args.assets));
const sliceSize = Math.ceil(args.assets / workerCount);
t0 = process.hrtime.bigint();
const promises = [];
for (let w = 0; w < workerCount; w++) {
const start = w * sliceSize;
const end = Math.min(start + sliceSize, args.assets);
if (start >= end) continue;
const panelSlice = panel.slice(start, end);
promises.push(
new Promise((resolve, reject) => {
const worker = new Worker(__filename, {
workerData: { panelSlice, indicator: args.indicator },
});
worker.on('message', (msg) => resolve({ start, results: msg }));
worker.on('error', reject);
worker.on('exit', (code) => {
if (code !== 0) reject(new Error(`worker exited with code ${code}`));
});
}),
);
}
Promise.all(promises)
.then((chunks) => {
const parallel = new Array(args.assets);
for (const { start, results } of chunks) {
for (let j = 0; j < results.length; j++) parallel[start + j] = results[j];
}
const tParallel = Number(process.hrtime.bigint() - t0) / 1e9;
const speedup = tSerial / Math.max(tParallel, 1e-9);
console.log(
`Parallel: ${tParallel.toFixed(3).padStart(8)} s ` +
`(workers=${workerCount}, speedup ~${speedup.toFixed(2)}x)`,
);
for (let i = 0; i < args.assets; i++) {
const s = serial[i];
const p = parallel[i];
if (s !== p && !(s === null && p === null)) {
console.error(`mismatch at asset ${i}: serial=${s} parallel=${p}`);
process.exit(1);
}
}
console.log('Parallel results match serial results — OK.');
})
.catch((err) => {
console.error(`parallel run failed: ${err.message}`);
process.exit(1);
});
}