jsrosetta

Bất đồng bộ

Structured concurrency: Promise.all và errgroup

Promise.all/allSettled/any kết hợp AbortController của Node.js so với errgroup, sync.WaitGroup và channel của Go: huỷ các tác vụ anh em khi một tác vụ lỗi, giới hạn số tác vụ chạy cùng lúc, gom mọi kết quả và lấy kết quả nhanh nhất.

Tác giả:
Phiên bản tối thiểu
Node.js ≥ 17.2Go ≥ 1.25
Đã chạy thử trên
Node.js 24.12.0Go 1.27.1

Code Node.js là ES module: lưu file .mjs hoặc đặt "type": "module" trong package.json.

"Structured concurrency" là một nguyên tắc đơn giản: mọi tác vụ con được khởi động trong một phạm vi phải kết thúc trước khi phạm vi đó kết thúc — không có tác vụ nào "trôi" ra ngoài, và lỗi của một tác vụ con được báo về cho phạm vi cha. Node.js có sẵn các bộ kết hợp Promise.all/allSettled/any/race, nhưng chúng chỉ chờ — không huỷ gì cả: khi Promise.all reject vì một promise lỗi, các promise còn lại vẫn chạy tiếp cho tới xong. Muốn các tác vụ anh em dừng lại, phải tự truyền vào một AbortSignal (xem bài huỷ tác vụ). Go không có kiểu Promise; các mẫu tương đương được ghép từ goroutine, sync.WaitGroup, channel và context — với errgroup (gói golang.org/x/sync, nằm ngoài thư viện chuẩn) là mảnh ghép gần với "Promise.all có huỷ" nhất.

Lỗi đầu tiên huỷ các tác vụ anh em (Promise.all + AbortController / errgroup)

import { setTimeout as sleep } from "node:timers/promises";
 
// work giả lập một việc mất `ms` mili giây, có thể thất bại, và dừng sớm
// khi signal bị abort.
async function work({ name, ms, fail }, signal) {
  await sleep(ms, undefined, { signal });
  if (fail) throw new Error(`${name} failed`);
  return `${name} ok`;
}
 
// all: như Promise.all, nhưng lỗi đầu tiên abort các task anh em, và chỉ
// trả về khi MỌI task đã dừng hẳn (giống errgroup.Wait). Trạng thái của task i
// được ghi vào statuses[i], như mảng statuses mà main bên Go cấp sẵn.
async function all(tasks, statuses) {
  const controller = new AbortController();
  const settled = await Promise.allSettled(
    tasks.map(async (task, i) => {
      try {
        const value = await work(task, controller.signal);
        statuses[i] = "ok";
        return value;
      } catch (err) {
        statuses[i] = err.name === "AbortError" ? "cancelled" : "failed";
        controller.abort(err); // abort lần hai trở đi không làm gì: reason giữ lỗi đầu tiên
        throw err;
      }
    }),
  );
  if (controller.signal.aborted) throw controller.signal.reason;
  return settled.map((s) => s.value);
}
 
const tasks = [
  { name: "A", ms: 100 },
  { name: "B", ms: 50, fail: true }, // B lỗi ở ~50ms, trước khi A và C xong
  { name: "C", ms: 300 },
];
const statuses = new Array(tasks.length);
 
try {
  await all(tasks, statuses); // chờ MỌI task dừng hẳn, rồi ném lỗi đầu tiên
} catch (err) {
  console.log(statuses); // → [ 'cancelled', 'failed', 'cancelled' ]
  console.log("error:", err.message); // → error: B failed
}

Giới hạn số tác vụ chạy cùng lúc (mapLimit / g.SetLimit)

import { setTimeout as sleep } from "node:timers/promises";
 
// mapLimit: chạy fn trên mọi item nhưng không quá `limit` việc cùng lúc —
// Node.js không có sẵn, nên tự dựng bằng `limit` "worker" cùng lấy từ một hàng đợi.
async function mapLimit(items, limit, fn) {
  const results = new Array(items.length);
  let next = 0;
  let failed = false;
 
  async function worker() {
    while (!failed && next < items.length) {
      const i = next++; // an toàn: JS đơn luồng, không có gì chen vào giữa đọc và tăng
      try {
        results[i] = await fn(items[i]);
      } catch (err) {
        failed = true; // worker khác thấy cờ này và ngừng lấy item mới
        throw err;
      }
    }
  }
 
  await Promise.all(Array.from({ length: limit }, () => worker()));
  return results;
}
 
let inFlight = 0;
let maxInFlight = 0;
 
const squares = await mapLimit([1, 2, 3, 4, 5, 6], 2, async (n) => {
  inFlight++;
  maxInFlight = Math.max(maxInFlight, inFlight);
  await sleep(20); // giả lập I/O
  inFlight--;
  return n * n;
});
 
console.log(squares);     // → [ 1, 4, 9, 16, 25, 36 ]
console.log(maxInFlight); // → 2 (tối đa)

Gom mọi kết quả, kể cả lỗi (Promise.allSettled / WaitGroup + errors.Join)

import { setTimeout as sleep } from "node:timers/promises";
 
async function check(name, ms, fail = false) {
  await sleep(ms);
  if (fail) throw new Error(`${name} unreachable`);
  return `${name} healthy`;
}
 
// allSettled không bao giờ reject: chờ mọi promise, kể cả khi vài cái lỗi
const settled = await Promise.allSettled([
  check("db", 50),
  check("cache", 20, true),
  check("queue", 30, true),
]);
 
for (const s of settled) {
  console.log(s.status, s.status === "fulfilled" ? s.value : s.reason.message);
}
// → fulfilled db healthy
// → rejected cache unreachable
// → rejected queue unreachable

Lấy kết quả thành công đầu tiên (Promise.any / channel + cancel)

import { setTimeout as sleep } from "node:timers/promises";
 
async function fetchFrom({ name, ms, fail }, signal) {
  await sleep(ms, undefined, { signal });
  if (fail) throw new Error(`${name} down`);
  return `data from ${name}`;
}
 
async function fastest(mirrors) {
  const controller = new AbortController();
  try {
    // any: lấy kết quả THÀNH CÔNG đầu tiên, bỏ qua các lần thất bại
    return await Promise.any(mirrors.map((m) => fetchFrom(m, controller.signal)));
  } finally {
    controller.abort(); // huỷ các request còn đang chạy — Promise.any không tự làm việc này
  }
}
 
console.log(await fastest([
  { name: "eu", ms: 300 },
  { name: "us", ms: 100 },
  { name: "asia", ms: 50, fail: true }, // lỗi sớm nhất, nhưng any bỏ qua
])); // → data from us
 
try {
  await fastest([{ name: "eu", ms: 30, fail: true }, { name: "us", ms: 10, fail: true }]);
} catch (err) {
  console.log(err.name, err.errors.map((e) => e.message)); // → AggregateError [ 'eu down', 'us down' ]
}

Khác biệt chính

Node.js Go
Lỗi đầu tiên reject ngay (không chờ, không huỷ) Promise.all không có tương đương trực tiếp
Chờ tất cả rồi trả lỗi đầu tiên Promise.allSettled + ném lại lỗi (hàm all ở trên) errgroup.Group.Wait()
Lỗi đầu tiên huỷ anh em AbortController truyền thủ công errgroup.WithContext
Giới hạn số tác vụ cùng lúc tự viết (hoặc thư viện như p-limit) g.SetLimit(n)
Gom mọi kết quả và lỗi Promise.allSettled sync.WaitGroup + errors.Join
Thành công đầu tiên Promise.any → AggregateError buffered channel + cancel() → errors.Join
Settle đầu tiên Promise.race giá trị đầu tiên nhận từ channel