如何在fp-ts中限制Task并发?新手团队求指导
嘿,作为刚上手fp-ts时也纠结过并发控制的人,我太懂你说的分块执行的痛点了——单个慢任务会把整个块的进度卡住,完全达不到真正的“并发上限”效果。其实fp-ts本身确实没有像Bluebird那样直接提供带concurrency参数的map,但我们可以利用Task的特性封装出符合FP风格的并发控制工具,而且思路比分块要灵活得多。
为什么分块不是最优解?
你提到的分块+sequence的方式,本质是“批量串行”——必须等当前块的所有任务都完成,才会启动下一块。如果块里有一个耗时特别长的任务,哪怕其他任务都快做完了,下一块的任务也得等着,整体效率会打折扣。我们需要的是“始终保持N个任务在运行”的真正并发限制:任务完成一个就补一个,不会因为单个慢任务阻塞后续任务的启动。
惯用实现:基于队列的并发控制函数
我们可以封装一个taskMapWithConcurrency函数,核心思路是维护一个任务队列,始终保持不超过指定数量的任务在执行。这种方式既符合fp-ts的Task抽象,又能达到和BluebirdPromise.map({concurrency: N})一样的效果。
代码实现
首先导入必要的fp-ts模块:
import * as T from 'fp-ts/lib/Task' import { pipe } from 'fp-ts/lib/function'
然后实现并发控制函数:
const taskMapWithConcurrency = <A>(concurrency: number) => (tasks: Array<T.Task<A>>): T.Task<Array<A>> => { // 边界情况处理 if (concurrency <= 0) throw new Error('Concurrency must be a positive number') if (tasks.length === 0) return T.of([]) return new T.Task((resolve) => { let completedTasks = 0 let nextTaskIndex = 0 const results: Array<A> = new Array(tasks.length) // 保持原任务顺序的结果数组 // 定义启动下一个任务的逻辑 const runNextTask = () => { // 所有任务都已启动,等待全部完成 if (nextTaskIndex >= tasks.length) { if (completedTasks === tasks.length) resolve(results) return } const currentIndex = nextTaskIndex nextTaskIndex++ // 执行当前任务,完成后更新结果并启动下一个 tasks[currentIndex]().then((result) => { results[currentIndex] = result completedTasks++ runNextTask() }) } // 启动初始的N个并发任务 const initialTasks = Math.min(concurrency, tasks.length) for (let i = 0; i < initialTasks; i++) { runNextTask() } }) }
使用示例
我们模拟几个不同耗时的任务来测试:
// 创建一个带延迟的测试任务 const createDelayedTask = (id: number, delayMs: number): T.Task<number> => T.of(id).chain(() => new T.Task((resolve) => setTimeout(() => resolve(id), delayMs))) // 任务数组:有快有慢 const tasks = [ createDelayedTask(1, 3000), // 最慢的任务 createDelayedTask(2, 1000), createDelayedTask(3, 2000), createDelayedTask(4, 500), createDelayedTask(5, 1500), ] // 限制并发数为2,执行任务并打印结果 pipe( tasks, taskMapWithConcurrency(2), T.map((results) => console.log('最终结果(保持原顺序):', results)) // 输出: [1,2,3,4,5] )()
这个示例中,任务1会运行3秒,但在它运行的过程中,任务2、4会先后完成,然后任务3、5会依次启动——始终保持最多2个任务在运行,不会因为任务1的慢而卡住后续任务的启动。
扩展到TaskEither(错误处理场景)
如果你的场景需要处理错误(比如API请求),只需要把Task换成TaskEither,修改任务完成后的逻辑即可:
import * as TE from 'fp-ts/lib/TaskEither' const taskEitherMapWithConcurrency = <E, A>(concurrency: number) => (tasks: Array<TE.TaskEither<E, A>>): TE.TaskEither<E, Array<A>> => { if (concurrency <= 0) return TE.left(new Error('Concurrency must be positive')) if (tasks.length === 0) return TE.right([]) return new TE.TaskEither((resolve) => { let completed = 0 let index = 0 const results: Array<A> = new Array(tasks.length) let hasError: E | null = null const runNext = () => { if (hasError || index >= tasks.length) { if (completed === tasks.length) { hasError ? resolve(TE.left(hasError)) : resolve(TE.right(results)) } return } const currentIndex = index index++ tasks[currentIndex]().then((either) => { if (either._tag === 'Left') { hasError = either.left resolve(TE.left(either.left)) } else { results[currentIndex] = either.right completed++ runNext() } }) } for (let i = 0; i < Math.min(concurrency, tasks.length); i++) { runNext() } }) }
总结
fp-ts的设计哲学是提供基础抽象,让开发者组合出需要的功能。这种基于队列的并发控制函数完全符合FP的思想,而且可以复用在各种异步场景中。相比分块执行,它能更高效地利用资源,不会被单个慢任务拖慢整体流程。
内容的提问来源于stack exchange,提问作者Tim

