| name | bun-workers |
| description | Use for Web Workers in Bun, worker_threads, parallel processing, and background tasks. |
| metadata | {"version":"1.0.0"} |
| license | MIT |
Bun Workers
Bun supports Web Workers and Node.js worker_threads for parallel execution.
Web Workers
Basic Usage
const worker = new Worker(new URL("./worker.ts", import.meta.url));
worker.postMessage({ type: "start", data: [1, 2, 3, 4, 5] });
worker.onmessage = (event) => {
console.log("Result:", event.data);
};
worker.onerror = (error) => {
console.error("Worker error:", error.message);
};
self.onmessage = (event) => {
const { type, data } = event.data;
if (type === "start") {
const result = data.map((x) => x * 2);
self.postMessage(result);
}
};
Worker with URL
const worker = new Worker(new URL("./worker.ts", import.meta.url));
const code = `
self.onmessage = (e) => {
self.postMessage(e.data * 2);
};
`;
const blob = new Blob([code], { type: "application/javascript" });
const worker = new Worker(URL.createObjectURL(blob));
Transferable Objects
const buffer = new ArrayBuffer(1024 * 1024);
const view = new Uint8Array(buffer);
view.fill(42);
worker.postMessage({ buffer }, [buffer]);
self.onmessage = (event) => {
const { buffer } = event.data;
const view = new Uint8Array(buffer);
self.postMessage({ buffer }, [buffer]);
};
Shared Memory
const shared = new SharedArrayBuffer(1024);
const view = new Int32Array(shared);
worker.postMessage({ shared });
Atomics.add(view, 0, 1);
self.onmessage = (event) => {
const { shared } = event.data;
const view = new Int32Array(shared);
Atomics.add(view, 0, 1);
Atomics.notify(view, 0);
};
Node.js worker_threads
import { Worker, isMainThread, parentPort, workerData } from "worker_threads";
if (isMainThread) {
const worker = new Worker(import.meta.filename, {
workerData: { numbers: [1, 2, 3, 4, 5] },
});
worker.on("message", (result) => {
console.log("Result:", result);
});
worker.on("error", (err) => {
console.error("Error:", err);
});
worker.on("exit", (code) => {
console.log("Worker exited with code:", code);
});
} else {
const { numbers } = workerData;
const sum = numbers.reduce((a, b) => a + b, 0);
parentPort?.postMessage(sum);
}
Worker Pool
import { Worker } from "worker_threads";
class WorkerPool {
private workers: Worker[] = [];
private queue: Array<{
task: any;
resolve: (value: any) => void;
reject: (err: Error) => void;
}> = [];
private activeWorkers = new Set<Worker>();
constructor(
private workerPath: string,
private poolSize: number
) {
for (let i = 0; i < poolSize; i++) {
this.addWorker();
}
}
private addWorker() {
const worker = new Worker(this.workerPath);
worker.on("message", {
..(worker);
.();
});
worker.(, {
..(worker);
.(, err);
});
..(worker);
}
(: ): <> {
( {
..({ task, resolve, reject });
.();
});
}
() {
( worker .) {
(!..(worker) && .. > ) {
{ task, resolve, reject } = ..()!;
..(worker);
worker.(, resolve);
worker.(, reject);
worker.(task);
}
}
}
() {
..( w.());
}
}
pool = (, );
results = .([
pool.({ : }),
pool.({ : }),
pool.({ : }),
]);
pool.();
Patterns
CPU-Intensive Tasks
const worker = new Worker(new URL("./cpu-worker.ts", import.meta.url));
const data = Array.from({ length: 1000000 }, () => Math.random());
worker.postMessage({ type: "process", data });
worker.onmessage = (event) => {
if (event.data.type === "progress") {
console.log(`Progress: ${event.data.percent}%`);
} else if (event.data.type === "result") {
console.log("Done:", event.data.result);
}
};
self.onmessage = (event) => {
const { type, data } = event.data;
if (type === "process") {
chunkSize = ;
result = ;
( i = ; i < data.; i++) {
result += .(data[i]);
(i % chunkSize === ) {
self.({
: ,
: .((i / data.) * ),
});
}
}
self.({ : , result });
}
};
Parallel Map
async function parallelMap<T, R>(
items: T[],
fn: string,
workerUrl: URL,
concurrency = 4
): Promise<R[]> {
const results: R[] = new Array(items.length);
const workers: Worker[] = [];
for (let i = 0; i < concurrency; i++) {
workers.push(new Worker(workerUrl));
}
let nextIndex = 0;
const processNext = (worker: Worker): Promise<void> => {
return new Promise((resolve) => {
if (nextIndex >= items.length) {
resolve();
return;
}
const index = nextIndex++;
worker.postMessage({ fn, item: items[index], index });
worker.onmessage = (event) => {
results[event.data.index] = event..;
(worker).(resolve);
};
});
};
.(workers.(processNext));
workers.( w.());
results;
}
Message Channel
const channel = new MessageChannel();
const worker1 = new Worker(new URL("./worker1.ts", import.meta.url));
const worker2 = new Worker(new URL("./worker2.ts", import.meta.url));
worker1.postMessage({ port: channel.port1 }, [channel.port1]);
worker2.postMessage({ port: channel.port2 }, [channel.port2]);
let port: MessagePort;
self.onmessage = (event) => {
if (event.data.port) {
port = event.data.port;
port.onmessage = (e) => console.log("From worker2:", e.data);
port.postMessage("Hello from worker1!");
}
};
Error Handling
const worker = new Worker(new URL("./worker.ts", import.meta.url));
worker.onerror = (error) => {
console.error("Uncaught error in worker:", error.message);
error.preventDefault();
};
worker.onmessageerror = (event) => {
console.error("Message deserialization failed");
};
self.onerror = (error) => {
self.postMessage({ type: "error", message: error.message });
};
Termination
const worker = new Worker(new URL("./worker.ts", import.meta.url));
worker.postMessage({ type: "shutdown" });
setTimeout(() => {
worker.terminate();
}, 5000);
self.onmessage = (event) => {
if (event.data.type === "shutdown") {
self.close();
}
};
Common Errors
| Error | Cause | Fix |
|---|
Worker not found | Wrong URL | Check worker file path |
Cannot serialize | Non-transferable data | Use transferable objects |
DataCloneError | Functions/DOM in message | Send only serializable data |
Worker terminated | Premature terminate | Check termination logic |
When to Load References
Load references/optimization.md when:
- Worker pool tuning
- Memory management
- Performance profiling
Load references/patterns.md when:
- Complex coordination
- Backpressure handling
- Error recovery