Chuyển đến nội dung chính

レッスン 9: ワーカー スレッドと CPU 負荷の高いタスク

ワーカー スレッド、SharedArrayBuffer、アトミックス。 MessageChannel、転送可能なオブジェクト。スレッド プール パターン、Piscina。 CPU 負荷の高いタスクのオフロード、画像処理。

💻 プログラミング — レッスン 9 レッスン 9: ワーカー スレッドと CPU 集中型 タスク

Node.js コア: 基本から高度まで

パート 3: 同時実行性とネットワーキング

xdev.asia

1. 基本的なワーカー スレッド

// main.ts
import { Worker, isMainThread, parentPort, workerData } from 'node:worker_threads'

if (isMainThread) {
  // Main thread — tạo worker
  const worker = new Worker(new URL(import.meta.url), {
    workerData: { start: 0, end: 1_000_000 }
  })

  worker.on('message', (result) => {
    console.log(`Sum: ${result}`)
  })

  worker.on('error', (err) => console.error(err))
  worker.on('exit', (code) => console.log(`Worker exited: ${code}`))
} else {
  // Worker thread
  const { start, end } = workerData as { start: number; end: number }
  let sum = 0
  for (let i = start; i < end; i++) sum += i
  parentPort!.postMessage(sum)
}

2. 並列処理

import { Worker } from 'node:worker_threads'
import os from 'node:os'

function runWorker(data: { start: number; end: number }): Promise<number> {
  return new Promise((resolve, reject) => {
    const worker = new Worker('./worker.js', { workerData: data })
    worker.on('message', resolve)
    worker.on('error', reject)
  })
}

// Chia công việc cho nhiều workers
const cpuCount = os.cpus().length
const total = 10_000_000
const chunkSize = Math.ceil(total / cpuCount)

const tasks = Array.from({ length: cpuCount }, (_, i) => ({
  start: i * chunkSize,
  end: Math.min((i + 1) * chunkSize, total),
}))

const results = await Promise.all(tasks.map(runWorker))
const totalSum = results.reduce((a, b) => a + b, 0)
console.log(`Total: ${totalSum} (${cpuCount} workers)`)

3. SharedArrayBuffer とアトミックス

// Shared memory giữa threads — zero-copy
const sharedBuffer = new SharedArrayBuffer(4) // 4 bytes
const sharedArray = new Int32Array(sharedBuffer)

const worker = new Worker('./worker.js', {
  workerData: { buffer: sharedBuffer }
})

// worker.js
const { buffer } = workerData
const arr = new Int32Array(buffer)

// Atomic operations — thread-safe
Atomics.add(arr, 0, 1)
Atomics.store(arr, 0, 42)
const value = Atomics.load(arr, 0)

// Wait/Notify — synchronization
Atomics.wait(arr, 0, 0)    // Block until value changes
Atomics.notify(arr, 0, 1)  // Wake 1 waiting thread

4. Piscina を使用したスレッド プール

import Piscina from 'piscina'

const pool = new Piscina({
  filename: new URL('./worker.js', import.meta.url).href,
  maxThreads: 4,
  minThreads: 2,
  idleTimeout: 30000,
})

// Worker task
// worker.js
export default function processImage(data: { path: string; width: number }) {
  // CPU-intensive image processing
  return { processed: true, path: data.path }
}

// Main — submit tasks
const results = await Promise.all(
  images.map(img => pool.run({ path: img, width: 800 }))
)

5. 譲渡可能なオブジェクト

// Transfer ownership (zero-copy) thay vì clone
const buffer = new ArrayBuffer(1024 * 1024) // 1MB

worker.postMessage({ buffer }, [buffer])
// buffer.byteLength === 0 (đã transfer, không còn ở main thread)

次の記事: クラスターモジュールとロードバランシング — フォーク、PM2、正常なシャットダウン。