redux-saga 通道(Channels)完全指南:actionChannel / eventChannel / channel / multicastChannel 详解
发布时间:2026/9/20 23:03:34来源:尧图网络
redux-saga 通道Channels完全指南actionChannel / eventChannel / channel / multicastChannel 详解【免费下载链接】redux-sagaAn alternative side effect model for Redux apps项目地址: https://gitcode.com/gh_mirrors/re/redux-saga通道Channels是 redux-saga 中连接外部事件源、在多个 Saga 之间传递消息的核心抽象。它把take/put从「与 Redux Store 通信」推广到「与任意事件源或 Saga 之间通信」同时支持消息缓冲是控制并发、串行化任务、实现 Worker 负载均衡的底层基础设施。读完本文你将掌握actionChannel、eventChannel、channel、multicastChannel四种通道的完整用法、缓冲策略选择以及它们在实际项目如 WebSocket 订阅、倒计时、并发受限任务池中的落地模式。本文以官方文档 docs/advanced/Channels.md 为核心骨架结合仓库源码packages/core/src/internal/channel.js、packages/core/src/internal/buffers.js 等与示例工程examples/cancellable-counter/src/sagas/index.js进行纵深展开。为什么需要 Channel从 watch-and-fork 到可控并发在此之前我们一直用take和put两个 Effect 与 Redux Store 通信。通道Channels是对这两类 Effect 的泛化它们既可以连接外部事件源也可以在 Saga 之间相互通信还可以对来自 Store 的特定 action 进行排队缓冲。先看经典的watch-and-fork模式import { take, fork, ... } from redux-saga/effects function* watchRequests() { while (true) { const {payload} yield take(REQUEST) yield fork(handleRequest, payload) } } function* handleRequest(payload) { ... }saga用fork启动非阻塞任务避免因阻塞而漏掉 Store 中的任何 action——每个REQUESTaction 都会创建一个handleRequest任务。如果 action 以极快速度大量触发就会同时存在大量并发执行的handleRequest任务并发度不受任何限制。现在假设需求变了我们希望串行处理REQUEST。任何时刻若有 4 个 action应先处理第一个REQUEST处理完成后再处理第二个依此类推。也就是说我们需要把所有尚未处理的 action 排队当前请求处理完毕后再从队列中取下一个消息。这正好是通道擅长的场景。本文后续四节将分别讲解四类通道通道形态数据来源典型用途actionChannel(pattern, [buffer])EffectRedux Store按 pattern 过滤的 action将特定 action 排队串行/限流处理eventChannel(subscriber, [buffer])工厂函数非 Effect任意外部事件源定时器、WebSocket 等把外部事件接入 Saga 的takechannel([buffer])工厂函数手动put的消息Saga 之间通信、Worker 负载均衡multicastChannel()工厂函数手动put的消息广播消息给多个不同 Worker使用actionChannelEffect为 Store 的 action 排队actionChannel是 redux-saga 提供的一个辅助 Effect可以帮我们实现「排队 串行处理」。重写上面的例子import { take, actionChannel, call, ... } from redux-saga/effects function* watchRequests() { // 1- Create a channel for request actions const requestChan yield actionChannel(REQUEST) while (true) { // 2- take from the channel const {payload} yield take(requestChan) // 3- Note that were using a blocking call yield call(handleRequest, payload) } } function* handleRequest(payload) { ... }三步走创建 action channelyield actionChannel(pattern)其中pattern的解释规则与之前take(pattern)完全一致字符串按action.type pattern精确匹配、数组匹配多个类型、函数作为谓词、symbol 按类型匹配详见 docs/API.md 中的take(pattern)一节。与take(pattern)的关键区别是actionChannel可以在 Saga 尚未准备好 take 时例如阻塞在一个 API 调用上缓冲 incoming 消息。从通道 takeyield take(requestChan)。take除了按 pattern 从 Store 取 action也可以接收通道对象。此时take会阻塞 Saga直到通道上有消息可用如果底层缓冲区已有积压消息take会立即恢复。使用阻塞式call这是串行化的关键。Saga 会一直阻塞到call(handleRequest)返回。在此期间若又有REQUESTaction 被 dispatch它们会被requestChan内部排队。当 Saga 从call(handleRequest)恢复、执行下一次yield take(requestChan)时take 会直接拿到排队的消息。底层实现视角actionChannel 如何拦截 Store action从源码看actionChannel(pattern, buffer)只是创建一个描述对象effect真正的工作在 packages/core/src/internal/io.js 的actionChannel与 packages/core/src/internal/effectRunnerMap.js 的runChannelEffect中完成// packages/core/src/internal/effectRunnerMap.js function runChannelEffect(env, { pattern, buffer }, cb) { const chan channel(buffer) // 创建底层通道默认 expanding 缓冲 const match matcher(pattern) // 把 pattern 编译为匹配函数 const taker (action) { if (!isEnd(action)) { env.channel.take(taker, match) // 继续监听 Store 上的匹配 action } chan.put(action) // 转发进内部通道供 Saga take } ... env.channel.take(taker, match) // 在 stdChannel 上注册监听 cb(chan) }也就是说actionChannel会在 Redux 的stdChannelmiddleware 内部通道上注册一个持续的监听者把所有匹配的 action转发进一个内部channel而 Saga 从这个内部通道取消息。actionChannel描述对象在 packages/core/src/internal/io.js 中还会校验pattern与buffer的合法性// packages/core/src/internal/io.js export function actionChannel(pattern, buffer) { check(pattern, is.pattern, actionChannel(pattern,...): argument pattern is not valid) if (arguments.length 1) { check(buffer, is.notUndef, actionChannel(pattern, buffer): argument buffer is undefined) check(buffer, is.buffer, actionChannel(pattern, buffer): argument ${buffer} is not a valid buffer) } return makeEffect(effectTypes.ACTION_CHANNEL, { pattern, buffer }) }控制缓冲默认无上限可自定义 Buffer默认情况下actionChannel无限制地缓冲所有 incoming 消息。如果你希望对缓冲做更多控制可以为 Effect 创建器提供 Buffer 参数。redux-saga 提供了几种常用缓冲none、dropping、sliding你也可以提供自己的 Buffer 实现Buffer 需要满足isEmpty()/put()/take()/flush()接口可参考 packages/core/src/internal/buffers.js 中ringBuffer的实现。例如只处理最近的 5 条消息最旧的消息会被丢弃import { buffers } from redux-saga import { actionChannel } from redux-saga/effects function* watchRequests() { const requestChan yield actionChannel(REQUEST, buffers.sliding(5)) ... }仓库提供的全部缓冲策略见 docs/API.md 的buffers一节实现在 packages/core/src/internal/buffers.js缓冲工厂溢出行为底层实现要点buffers.none()无缓冲没有 pending taker 时新消息直接丢失空实现zeroBufferbuffers.fixed(limit)缓冲到limit条溢出抛Error省略limit默认 10ringBuffer(limit, ON_OVERFLOW_THROW)buffers.expanding(initialSize)同fixed但溢出时容量动态翻倍扩容ringBuffer(initialSize, ON_OVERFLOW_EXPAND)buffers.dropping(limit)同fixed但溢出时静默丢弃新消息ringBuffer(limit, ON_OVERFLOW_DROP)buffers.sliding(limit)同fixed但溢出时新消息插入队尾、丢弃最旧消息ringBuffer(limit, ON_OVERFLOW_SLIDE)actionChannel的默认缓冲区是expanding见 packages/core/src/internal/channel.js 中channel(buffer buffers.expanding())的默认参数这也是「默认无上限缓冲」这一行为的事实来源。从 packages/core/src/internal/buffers.js 的ringBuffer实现可以看到几种溢出策略的具体差异fixed缓冲满时调用put会抛出Channels Buffer overflow!sliding满时arr[pushIndex] it直接覆盖当前位置再同步推进pushIndex与popIndex等效于丢掉最旧的一条expanding满时先把现有元素flush()出来limit翻倍后重新入队dropping满时什么都不做新消息被静默丢弃。使用eventChannel工厂连接外部事件源与actionChannelEffect不同eventChannel是一个工厂函数它同样创建通道但数据来源是Redux Store 之外的事件源。一个基础的倒计时示例——从 interval 创建通道import { eventChannel, END } from redux-saga function countdown(secs) { return eventChannel(emitter { const iv setInterval(() { secs - 1 if (secs 0) { emitter(secs) } else { // this causes the channel to close emitter(END) } }, 1000); // The subscriber must return an unsubscribe function return () { clearInterval(iv) } } ) }eventChannel的第一个参数是subscriber订阅者函数。它的职责是初始化外部事件源上面用setInterval通过调用传入的emitter把事件源的每个 incoming 事件路由进通道。上面例子中我们每秒调用一次emitter。注意必须对事件源做净化不要把 null 或 undefined 传进 event channel。传数字本身没问题但我们建议像组织 redux action 一样组织 event channel 的数据——用{ number }而不是裸number。这里还调用了emitter(END)用于通知通道的消费者通道已关闭之后不会再有任何消息。在 Saga 中消费 event channel下面的用法取自仓库的 cancellable-counter 示例完整实现见 examples/cancellable-counter/src/sagas/index.jsimport { take, put, call } from redux-saga/effects import { eventChannel, END } from redux-saga // creates an event Channel from an interval of seconds function countdown(seconds) { ... } export function* saga() { const chan yield call(countdown, value) try { while (true) { // take(END) will cause the saga to terminate by jumping to the finally block let seconds yield take(chan) console.log(countdown: ${seconds}) } } finally { console.log(countdown terminated) } }Saga 执行yield take(chan)后会阻塞直到通道上有消息——对应示例中调用emitter(secs)的时刻。注意我们是在try/finally中执行整个while (true) {...}循环interval 结束时countdown 函数通过emitter(END)关闭 event channel关闭通道会终止所有阻塞在该通道take上的 Saga。在我们的示例中终止 Saga 会使其跳到finally块如果提供了 finally否则 Saga 直接终止。从源码看eventChannel的关闭链路非常清晰packages/core/src/internal/channel.jsexport function eventChannel(subscribe, buffer buffers.none()) { ... unsubscribe subscribe((input) { if (isEnd(input)) { close(); return } // emitter(END) - 关闭通道 chan.put(input) }) ... return { take: chan.take, flush: chan.flush, close } }即emitter(END)会触发内部close()先调用unsubscribe()解除外部事件源订阅再关闭底层channel底层channel.close()会向所有 pending taker 投递END从而终止阻塞中的 Saga。注意这里一个关键事实eventChannel的默认缓冲是buffers.none()默认不缓冲消息。取消与提前退出chan.close()的正确姿势subscriber 返回一个unsubscribe函数。通道会在事件源完成之前用它在内部解绑。如果 Saga 想在事件源完成前提前退出例如 Saga 被取消可以调用chan.close()来关闭通道并从事件源解绑。给 Saga 加上取消支持import { take, put, call, cancelled } from redux-saga/effects import { eventChannel, END } from redux-saga // creates an event Channel from an interval of seconds function countdown(seconds) { ... } export function* saga() { const chan yield call(countdown, value) try { while (true) { let seconds yield take(chan) console.log(countdown: ${seconds}) } } finally { if (yield cancelled()) { chan.close() console.log(countdown cancelled) } } }cancellable-counter 示例examples/cancellable-counter/src/sagas/index.js中还展示了更完整的组合用race([call(incrementAsync, action), take(CANCEL_INCREMENT_ASYNC)])让「取消 action」赢得竞速后自动取消倒计时 Saga被取消的 Saga 在finally中通过cancelled()判断后执行chan.close()清理 interval。仓库对应的测试位于 examples/cancellable-counter/test/sagas.js。实战用 eventChannel 把 WebSocket 事件接入 Saga再看一个把 WebSocket 事件例如基于 socket.io接入 Saga 的完整示例等待服务器消息ping延迟一段时间后回复pong。import { take, put, call, apply, delay } from redux-saga/effects import { eventChannel } from redux-saga import { createWebSocketConnection } from ./socketConnection // this function creates an event channel from a given socket // Setup subscription to incoming ping events function createSocketChannel(socket) { // eventChannel takes a subscriber function // the subscriber function takes an emit argument to put messages onto the channel return eventChannel(emit { const pingHandler (event) { // puts event payload into the channel // this allows a Saga to take this payload from the returned channel emit(event.payload) } const errorHandler (errorEvent) { // create an Error object and put it into the channel emit(new Error(errorEvent.reason)) } // setup the subscription socket.on(ping, pingHandler) socket.on(error, errorHandler) // the subscriber must return an unsubscribe function // this will be invoked when the saga calls channel.close method const unsubscribe () { socket.off(ping, pingHandler) } return unsubscribe }) } // reply with a pong message by invoking socket.emit(pong) function* pong(socket) { yield delay(5000) yield apply(socket, socket.emit, [pong]) // call emit as a method with socket as context } export function* watchOnPings() { const socket yield call(createWebSocketConnection) const socketChannel yield call(createSocketChannel, socket) while (true) { try { // An error from socketChannel will cause the saga jump to the catch block const payload yield take(socketChannel) yield put({ type: INCOMING_PONG_PAYLOAD, payload }) yield fork(pong, socket) } catch(err) { console.error(socket error:, err) // socketChannel is still open in catch block // if we want end the socketChannel, we need close it explicitly // socketChannel.close() } } }关键细节Error 对象可以直接 emit 进通道。当 socket 出错时emit(new Error(...))会让阻塞在yield take(socketChannel)上的 Saga 抛错、跳入catch块。这一点对应 packages/core/src/internal/effectRunnerMap.js 中runTakeEffect的实现——takeCb会检查input instanceof Error并以错误回调结束。emit 错误不会默认关闭通道。catch 块执行后socketChannel仍然打开若想结束通道需显式调用socketChannel.close()。注意错误处理中只解绑了ping监听如需完整清理应在unsubscribe中同时移除error监听示例保持精简。注意eventChannel 上的消息默认不缓冲。如需指定缓冲策略要给 eventChannel 工厂传入 buffer例如eventChannel(subscriber, buffer)。缓冲实现与actionChannel完全一致buffers.none/fixed/expanding/dropping/sliding详见 docs/API.md 的buffers一节。使用channel在 Saga 之间通信Worker 池与负载均衡除了 action channel 和 event channel你还可以直接创建不连接任何数据源的通道然后手动put消息到通道上。当需要用通道在 Saga 之间通信时这非常方便。回到请求处理的例子import { take, fork, ... } from redux-saga/effects function* watchRequests() { while (true) { const {payload} yield take(REQUEST) yield fork(handleRequest, payload) } } function* handleRequest(payload) { ... }前面已经看到watch-and-fork 允许无限并发地同时处理多个请求随后我们用actionChannel把并发度限制为同一时刻 1 个任务。现在假设需求是同一时刻最多 3 个任务在跑。收到请求时如果正在执行的任务少于 3 个就立即处理否则把任务排队等待 3 个slot之一空出来。下面是用channel的解法import { channel } from redux-saga import { take, fork, ... } from redux-saga/effects function* watchRequests() { // create a channel to queue incoming requests const chan yield call(channel) // create 3 worker threads for (var i 0; i 3; i) { yield fork(handleRequest, chan) } while (true) { const {payload} yield take(REQUEST) yield put(chan, payload) } } function* handleRequest(chan) { while (true) { const payload yield take(chan) // process the request } }解析用channel工厂创建通道。默认情况下它会缓冲所有 put 进来的消息除非有 pending taker此时 taker 立即被消息恢复——这正是 packages/core/src/internal/channel.js 中channel()的逻辑put时若takers.length 0则写入 buffer否则直接调用第一个 taker 回调take时若 buffer 非空则立即取出否则把回调挂进takers队列等待。watchRequestsfork 出 3 个 worker saga同一个通道传给所有 fork 出的 saga。watchRequests用它向 3 个 worker「分发」工作每个REQUESTaction 到达就把 payloadput到通道。payload 会被任意一个空闲worker take 走否则由通道排队直到某个 worker Saga 准备好 take。3 个 worker 都跑典型的 while 循环每轮 take 下一个请求没有请求就阻塞等待。该机制在 3 个 worker 之间提供了自动负载均衡——快的 worker 不会被慢的 worker 拖慢。对应通道核心实现packages/core/src/internal/channel.jsfunction put(input) { if (closed) return if (takers.length 0) { return buffer.put(input) // 无等待者 - 入缓冲 } const cb takers.shift() // 有等待者 - 直接唤醒第一个 cb(input) } function take(cb) { if (closed buffer.isEmpty()) { cb(END) } else if (!buffer.isEmpty()) { cb(buffer.take()) // 缓冲有货 - 立即取出 } else { takers.push(cb) // 否则挂起等待 cb.cancel () { remove(takers, cb) } } }这正是「单播」语义一条消息只会被一个 taker 消费因此天然适配 Worker 池的任务分发。take的取消能力cb.cancel由 packages/core/src/internal/effectRunnerMap.js 的runTakeEffect透传给 Effect 取消流程例如race中输掉的分支会取消其挂起的 take。使用multicastChannel向不同Worker 广播上一节看到channel如何在同一个被 fork 多次的 worker之间做负载均衡。那如果需要put一个 action 到通道、让多个不同的 worker都消费它呢比如把同一个请求同时交给多个执行不同副作用side effect的 worker。先用普通channel演示问题yield put(chan, payload)永远只会唤醒一个workerlogWorker或mainWorker不会两个都执行import { channel } from redux-saga import { take, fork, call, put } from redux-saga/effects function* watchRequests() { // create a channel to queue incoming requests const chan yield call(channel) // fork both workers yield fork(logWorker, chan) yield fork(mainWorker, chan) while (true) { const { payload } yield take(REQUEST) // put here will reach only one worker, not both! yield put(chan, payload) } } function* logWorker(channel) { while (true) { const payload yield take(channel) // Log the request somewhere.. console.log(logWorker:, payload) } } function* mainWorker(channel) { while (true) { const payload yield take(channel) // Process the request console.log(mainWorker, payload) } }要解决这个问题需要使用multicastChannel它会把 action同时广播给所有 worker。注意对multicastChannel使用take时需要额外传入pattern参数——可以用它过滤要take的 action。import { multicastChannel } from redux-saga import { take, fork, call, put } from redux-saga/effects function* watchRequests() { // create a multicastChannel to queue incoming requests const channel yield call(multicastChannel) // fork different workers yield fork(logWorker, channel) yield fork(mainWorker, channel) while (true) { const { payload } yield take(REQUEST) yield put(channel, payload) } } function* logWorker(channel) { while (true) { // Pattern * for simplicity const payload yield take(channel, *) // Log the request somewhere.. console.log(logWorker:, payload) } } function* mainWorker(channel) { while (true) { // Pattern * for simplicity const payload yield take(channel, *) // Process the request console.log(mainWorker, payload) } }底层实现multicast 与 pattern 匹配从源码看packages/core/src/internal/channel.jsmulticastChannel与channel的关键差异在于export function multicastChannel() { ... return { [MULTICAST]: true, // 标记为多播通道 put(input) { if (closed) return if (isEnd(input)) { close(); return } const takers (currentTakers nextTakers) for (let i 0; i takers.length; i) { const taker takers[i] if (takerMATCH) { // 用 pattern 过滤 taker.cancel() // 每个 taker 只能消费一次 taker(input) } } }, take(cb, matcher matchers.wildcard) { ... cb[MATCH] matcher // 把 pattern 编译成的匹配函数挂在 taker 上 nextTakers.push(cb) ... }, close, } }要点put会遍历所有 taker对每个 taker 先做takerMATCHpattern 匹配命中才唤醒——这就是「广播」语义的实现每个 taker 被唤醒前先调用taker.cancel()把自己从等待列表移除保证每个 taker 每条消息只消费一次take(cb, matcher)的第二个参数默认matchers.wildcard见 packages/core/src/internal/matcher.js通配匹配恒为真。所以文档强调对 multicastChannel 使用take时建议显式传 pattern如*否则在 packages/core/src/internal/io.js 的take(patternOrChannel, multicastPattern)中仅当is.multicast(patternOrChannel) is.notUndef(multicastPattern) is.pattern(multicastPattern)时才会构造{ channel, pattern }形式——不传 pattern 会走普通take(channel)分支并打印警告take(channel) takes one argument but two were provided或忽略第二个参数。multicastChannel没有内部缓冲put 时只唤醒当前在等待的 taker没有等待者则消息直接丢弃这与channel默认缓冲的行为不同因此它不用于排队积压只用于即时广播。另外middleware 内部默认的stdChannel正是multicastChannel的增强版packages/core/src/internal/channel.js 的stdChannel()它把 Store dispatch 的 action 通过asap调度器异步投递SAGA_ACTION标记的 action 除外会同步投递。这也是take(pattern)能从 Store 匹配 action 的底层机制。注意multicastChannel没有消息缓冲也不提供队列语义因此它适合「一对多广播」若需要「一对多 排队」可以在广播前自行用channel/actionChannel做一层积压或由各 worker 自行处理。通道对比速查与选型建议维度actionChanneleventChannelchannelmulticastChannel类型Effect需yield工厂函数通常配合call工厂函数配合call工厂函数配合call数据来源Redux Store 中匹配 pattern 的 action任意外部事件源定时器、WebSocket 等手动put手动put默认缓冲expanding无上限none不缓冲expanding无上限无缓冲不排队消费语义单播一个消费者单播一个消费者单播一个消费者自动负载均衡多播广播给所有匹配 taker关闭方式通道关闭时停止转发 Store actionemitter(END)或chan.close()chan.close()chan.close()或put(END)典型场景串行化/限流处理 Store actionWebSocket、定时器等外部事件接入固定并发数的 Worker 池同一事件分发多个不同 Worker选型建议串行处理 Store action、限制处理速率→actionChannel(PATTERN, buffer) 阻塞call接入定时器、WebSocket、socket.io 等外部事件→eventChannel(subscriber, buffer)记得返回unsubscribe在finally中处理取消并chan.close()固定 N 个并发 worker、自动负载均衡→channel() fork N 个相同 worker同一份数据需要多个不同 worker 同时消费→multicastChannel()take(channel, pattern)。总结通道把 redux-saga 的take/put从「只能与 Redux Store 对话」扩展为「与任意事件源和任意 Saga 对话」的统一抽象actionChannel在 Store 与 Saga 之间加了一层可配置的缓冲队列是实现串行处理与限流的开箱即用方案eventChannel让外部事件源interval、WebSocket、socket.io 等以统一的消息形式进入 Saga 的世界配合END与chan.close()可以优雅处理完成与取消channel提供手动的单播消息队列天然支持固定并发 Worker 池与负载均衡multicastChannel提供一对多的广播语义配合take(channel, pattern)的过滤能力适合把同一事件分发给多个职责不同的 Worker。深入理解它们的差异后再回头看 redux-saga 内部middleware 默认的stdChannel本质就是一个异步投递的multicastChannel所有take(pattern)都建立在其上actionChannel又是通过在该通道上注册持续监听者、把匹配 action 转发进内部channel实现的见 packages/core/src/internal/effectRunnerMap.js。掌握这一层「通道即抽象」的心智模型你就拥有了在 Saga 架构中处理并发、排队与事件接入的完整工具箱。【免费下载链接】redux-sagaAn alternative side effect model for Redux apps项目地址: https://gitcode.com/gh_mirrors/re/redux-saga创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网