datapipeline类实现高并发异步任务调度的流式数据清洗管道,支持阶段注册、并发限制、分批处理、错误隔离及钩子扩展。

直接用 ES6 class 构建“高并发异步任务调度的流式数据清洗管道基类”,需要明确一点:JavaScript 单线程本质决定了它不提供原生的多线程并发能力,所谓“高并发”实际是高吞吐、低延迟、可并行调度的异步流水线,核心靠 Promise 链、任务队列控制、并发数限制(如 Promise.allSettled + 信号量)与可插拔的处理阶段设计。
1. 定义管道基类:支持阶段注册与顺序执行
基类不直接执行清洗逻辑,而是管理阶段(stages)、输入源(source)、输出目标(sink)和调度策略。每个 stage 是一个异步函数,接收数据并返回清洗后数据(或 reject 错误)。
关键设计点:
- 构造时接受可选的 maxConcurrency(默认 3),用于限制同时运行的 stage 实例数
- 用 Array.push() 累积 stage,保证执行顺序
- 所有 stage 必须返回 Promise,统一用 async/await 或 Promise.then 封装
示例代码:
class DataPipeline {constructor(maxConcurrency = 3) {
this.stages = [];
this.maxConcurrency = maxConcurrency;
}
use(stageFn) {
if (typeof stageFn !== 'function') throw new TypeError('Stage must be a function');
this.stages.push(stageFn);
return this;
}
}
2. 实现流式调度:按批+限流+错误隔离
清洗管道不能一次性 load 所有数据(内存溢出),也不应让一个失败 stage 阻塞整条流。推荐使用“分批处理 + 并发控制 + 失败跳过”模式。
核心方法 process(items) 应:
- 将输入数组切分为大小为 maxConcurrency 的批次
- 对每一批调用 Promise.allSettled(),确保单批内 stage 并行但互不干扰
- 每个 stage 调用包裹 try/catch,失败时不中断后续 stage,记录 error 或打标记
- 返回结构化结果:{ data: cleanedItems[], errors: [] }
示例片段:
async process(items) {const results = [];
const errors = [];
const batches = this.#chunk(items, this.maxConcurrency);
for (const batch of batches) {
const settled = await Promise.allSettled(
batch.map(item => this.#runStages(item))
);
settled.forEach(r => {
if (r.status === 'fulfilled') results.push(r.value);
else errors.push(r.reason);
});
}
return { data: results, errors };
}
#runStages(item) {
return this.stages.reduce((acc, stage) => acc.then(data => stage(data)), Promise.resolve(item));
}
3. 支持异步清洗阶段:每个 stage 可含 I/O 或计算
清洗阶段本身必须是异步友好的。比如去重查库、调用外部 API 校验手机号、格式化时间戳等。
正确写法(返回 Promise):
const validatePhone = async (record) => {const res = await fetch(`/api/validate?phone=${record.phone}`);
if (!res.ok) throw new Error(`API failed: ${res.status}`);
const valid = await res.json();
return { ...record, isValid: valid };
};
错误写法(同步阻塞、无 error 处理):
// ❌ 不要这样:const badStage = (item) => {
JSON.parse(item.raw); // 同步抛错会中断整个 pipeline
return item;
};
4. 扩展性设计:支持中间件式钩子与生命周期
真实清洗流程常需日志、指标上报、超时控制、重试等。可在基类中预留钩子:
- onStageStart(stageName, item):stage 开始前触发
- onStageError(stageName, item, error):stage 报错时触发
- onBatchComplete(batchIndex, resultCount):每批完成后触发
这些钩子默认为空函数,子类可 override 或通过 options 注入,不影响主流程。











