在 BEAM 上做数据处理时,最容易被忽略的一点是:数据的产生速度与消费速度几乎从不相等。一个从 Kafka 拉取的管道,可能在业务高峰期每秒涌入十万条消息,而下游写库只能处理三千条;如果中间没有机制让生产者「慢下来」,内存会在几十秒内被消息堆积撑爆,随后整个节点被 OOM Killer 杀死。GenStage 就是为了解决这个问题而生的抽象——它把「需求(demand)」变成流的第一性概念,让下游主动向上游索取数据,从而实现天然的背压。
本文从 GenStage 的需求驱动模型讲起,逐层展开三类角色的回调契约、Flow 的分区与窗口算子,最后落到 Broadway 这一生产级管道框架:如何在 SQS/Kafka/RabbitMQ 上编排并发、批处理、失败重试与死信队列,并接入 Telemetry 做端到端观测。
一、流式处理的挑战与 GenStage 定位
1.1 批处理与流处理的本质差异
| 维度 | 批处理 | 流处理 |
|---|---|---|
| 数据边界 | 有界(文件/表全量) | 无界(持续到达) |
| 触发方式 | 定时/手动调度 | 事件到达即触发 |
| 延迟 | 分钟到小时 | 毫秒到秒 |
| 状态管理 | 作业内临时状态 | 长期状态 + 窗口 |
| 失败恢复 | 整批重跑 | 逐条重试/死信 |
| 背压需求 | 弱(数据量已知) | 强(速率不可预测) |
流处理的核心难点不是「处理」,而是「在速率失配时如何优雅退化」。传统做法是给队列设一个上限,满了就阻塞或丢弃;而 GenStage 的思路是让需求自下而上流动,上游只在被索取时才生产。
1.2 GenStage 的由来
GenStage 由 José Valim 在 OTP 18 时代引入,作为 gen_event 的继任者。gen_event 是「推」模型——事件管理器不管订阅者是否跟得上,一股脑推送,慢订阅者只能靠进程邮箱堆积,最终导致内存膨胀甚至节点崩溃。GenStage 反转为「拉」模型:
生产者 ──事件──▶ 生产者消费者 ──事件──▶ 消费者
▲ ▲ │
└──── demand ──────┴──── demand ────┘
(需求自下而上流动)
每个环节只在上游有富余时推进,速率由最慢的一环决定,多余的数据留在源头(如 Kafka 的 offset 未提交、SQS 消息未被接收),而不是堆在内存里。
1.3 三类角色
| 角色 | 声明返回值 | 典型用途 |
|---|---|---|
| Producer | {:producer, state} | 数据源:Kafka 消费者、文件读取、轮询 API |
| Consumer | {:consumer, state} | 数据终点:写库、发 HTTP、落盘 |
| ProducerConsumer | {:producer_consumer, state} | 中间转换:解析、过滤、富化、分区 |
这三类角色组合起来就构成一张有向无环图,GenStage 的 GenStage.sync_subscribe/2 与 subscribe_to 负责在 init/1 阶段建立订阅关系并协商初始需求。
二、GenStage 核心概念与回调契约
2.1 最小生产者
defmodule CounterProducer do
use GenStage
def start_link(initial) do
GenStage.start_link(__MODULE__, initial, name: __MODULE__)
end
@impl true
def init(counter), do: {:producer, counter}
@impl true
def handle_demand(demand, counter) when demand > 0 do
# 只在被索取时才生成事件,demand 就是下游要的数量
events = Enum.to_list(counter..(counter + demand - 1))
{:noreply, events, counter + demand}
end
end
关键点:handle_demand/2 的返回值是 {:noreply, events, new_state},其中 events 的长度应当等于(或小于)demand。如果返回超过 demand 的事件数,GenStage 会发出警告——那意味着你在无节制地推数据。
2.2 最小消费者
defmodule PrintConsumer do
use GenStage
def start_link(_), do: GenStage.start_link(__MODULE__, :ok, name: __MODULE__)
@impl true
def init(:ok) do
# subscribe_to 在 init 中建立订阅,max_demand 决定一次最多索取多少
{:consumer, :ok, subscribe_to: [{CounterProducer, max_demand: 100}]}
end
@impl true
def handle_events(events, _from, state) do
Enum.each(events, &IO.inspect(&1, label: "event"))
{:noreply, [], state}
end
end
消费者处理完一批事件后,GenStage 会自动补发新的 demand,形成闭环。这个「自动补发」正是背压的全部秘密:处理慢 → 补发慢 → 上游生产慢。
2.3 生产者消费者与订阅选项
defmodule EnrichStage do
use GenStage
@impl true
def init(_) do
{:producer_consumer, %{},
subscribe_to: [{CounterProducer, max_demand: 500, min_demand: 250}]}
end
@impl true
def handle_events(events, _from, state) do
enriched = Enum.map(events, fn n -> %{value: n, enriched: n * 2} end)
{:noreply, enriched, state}
end
end
| 订阅选项 | 默认值 | 含义 |
|---|---|---|
max_demand | 1000 | 一次最多向上游索取的事件数 |
min_demand | div(max_demand, 2) | 缓冲区低于此值时触发补货 |
buffer_size | max_demand 的一半 | 缓冲上限,超出会丢弃 |
dispatcher | GenStage.DemandDispatcher | 多订阅者时的分配策略 |
2.4 需求记账公式
GenStage 内部维护一个需求计数器,语义可以简化为:
待处理量 = buffer_size + 已发出但未满足的 demand
补货条件 = 待处理量 < min_demand
补货数量 = max_demand - min_demand
理解这个公式就能解释绝大多数「为什么我的 GenStage 吞吐上不去」:如果 max_demand 设得太小(比如 1),每条消息都要走一次跨进程消息往返,吞吐被消息传递开销压死;设得太大(比如 100000),单批内存占用过高,GC 停顿变长。
三、需求驱动的背压机制
3.1 背压的传播路径
假设管道是「Kafka → 解析 → 写 PostgreSQL」,写库环节因为数据库锁竞争变慢:
写库变慢
└─▶ 消费者补发 demand 的频率下降
└─▶ 解析阶段停止向 Kafka 生产者索取
└─▶ Kafka 生产者不再 poll 新消息
└─▶ offset 不推进,消息留在 broker
整条链路的速率自动收敛到最慢环节,且没有任何一处需要显式限流配置。这是 GenStage 相比「手动加信号量」最大的优势。
3.2 缓冲区的双刃剑
buffer_size 允许上游在需求未被消费时先缓存一部分事件,起到削峰填谷的作用,例如 subscribe_to: [{MyProducer, max_demand: 1000, min_demand: 100, buffer_size: 1000}]。但缓冲区不是越大越好:
- 缓冲区内的数据不落盘,进程崩溃即丢失;
- 缓冲占用进程堆,堆越大 GC 的标记阶段越慢;
- 大量缓冲会掩盖下游的真实延迟,让「看起来正常」的系统在某一刻突然雪崩。
3.3 用 GenStage 做限流
需求驱动模型天然适合限流——只要把消费者的处理速度压下来,整条链就跟着慢下来。在 handle_events/3 里对每条事件调用 RateLimiter.acquire!(:upstream, 1)(令牌桶,每秒最多放行 N 条),上游就会因为 demand 得不到满足而停止生产。对比「在生产者端 sleep」,这种做法的好处是不浪费任何一次跨进程消息传递:上游根本没生产,而不是生产了再丢掉。
3.4 分区与并发
GenStage 的 DemandDispatcher 会把需求分发给多个订阅者,实现同一阶段的水平扩展:
# 启动 8 个消费者,共享同一生产者的需求
for i <- 1..8 do
GenStage.start_link(MyConsumer, :ok, name: :"consumer_#{i}")
end
默认 DemandDispatcher 采用轮询分配,适合事件大小均匀的场景;若事件处理代价差异极大,应改用 GenStage.PartitionDispatcher 按 key 分区,避免「某个消费者被重活拖死,其他消费者空转」。
四、Flow 分区、窗口与聚合
4.1 Flow 是什么
Flow 建立在 GenStage 之上,用函数式管道语法表达并行计算,并把「阶段划分」隐藏起来:
1..1_000_000
|> Flow.from_enumerable(max_demand: 1000)
|> Flow.map(&(&1 * 2))
|> Flow.filter(&(rem(&1, 3) == 0))
|> Flow.partition()
|> Flow.reduce(fn -> 0 end, fn _event, acc -> acc + 1 end)
|> Flow.emit(:state)
|> Enum.to_list()
Flow.from_enumerable/2 会自动把可枚举切成多块,分发到多个并行阶段;Flow.partition/2 是一个同步点,之后的所有算子都在分区内独立执行。
4.2 分区是并行聚合的前提
Flow 有一条铁律:Flow.reduce/3 之前必须 Flow.partition/2。原因是聚合需要状态,而状态不能跨阶段共享;分区把事件按 key 稳定地路由到同一台「逻辑机器」,聚合才有意义:
events
|> Flow.from_enumerable()
|> Flow.partition(key: {:key, :user_id}, stages: 4)
|> Flow.reduce(fn -> %{} end, fn event, acc ->
Map.update(acc, event.type, 1, &(&1 + 1))
end)
|> Flow.emit(:state)
|> Enum.to_list()
stages: 4 指定分区数量,超过分区数的并发没有意义——多余阶段会收到空分区。
4.3 窗口类型
无界流上的聚合必须限定时间或数量范围,Flow 提供三类窗口:
| 窗口 | 构造函数 | 触发条件 | 典型用途 |
|---|---|---|---|
| 计数窗口 | Flow.Window.count(1000) | 每满 1000 条 | 定批处理、分页上报 |
| 周期窗口 | Flow.Window.periodic(5, :second) | 每 5 秒 | 指标聚合、心跳统计 |
| 全局窗口 | Flow.Window.global() | 手动 Flow.emit_and_reduce/3 | 会话聚合、累计计数 |
Flow.from_enumerable(events)
|> Flow.partition(window: Flow.Window.count(1000), key: {:key, :tenant})
|> Flow.reduce(fn -> %{count: 0, bytes: 0} end, fn e, acc ->
%{acc | count: acc.count + 1, bytes: acc.bytes + byte_size(e.body)}
end)
|> Flow.emit(:state)
|> Enum.to_list()
4.4 触发器与延迟
窗口默认在关闭时触发一次。若需要「滑动」语义(每 10 条输出一次最近 100 条的统计),用 Flow.Window.count(100, trigger: Flow.Window.count(10)) 配置触发器;若需要等待乱序数据,则在周期窗口上加 latency,如 Flow.Window.periodic(5, :second, latency: 1),表示窗口在最后一个事件到达后再等 1 秒才关闭。
4.5 Flow 与 GenStage 的选型
| 场景 | 推荐 | 理由 |
|---|---|---|
| 固定数据集的并行转换 | Flow | 自动分块,代码最简 |
| 无界流 + 外部数据源 | Broadway | 内置生产者、ack、重试 |
| 需要精确控制拓扑 | 裸 GenStage | 完全掌控需求与调度 |
| 一次性 ETL 作业 | Flow | Flow.from_enumerable 直接吃文件/表 |
Flow 不适合长期驻留的服务:它没有内置的 ack/重试语义,进程崩溃后数据恢复要靠调用方保证。生产环境的常驻管道应当用 Broadway。
五、Broadway 生产管道
5.1 Broadway 的架构
Broadway 是官方维护的「GenStage 成品化封装」,把生产管道拆成四个可独立配置并发的阶段:
生产者(producer) ─▶ 处理器(processors) ─▶ 批处理器(batchers) ─▶ 消费者(consumers)
1~N 个 并发 M 个 并发 K 个 批处理器内部
| 阶段 | 职责 | 关键配置 |
|---|---|---|
| producer | 从外部系统拉取消息 | 模块 + 连接配置,通常并发 1 |
| processors | 逐条转换、调用业务逻辑 | concurrency、max_demand |
| batchers | 按 size/timeout 聚合 | batch_size、batch_timeout |
| consumers | 在 batcher 内部批量落地 | 与 batcher 同一进程 |
5.2 一个 SQS 管道
defmodule MyApp.SQSPipeline do
use Broadway
alias Broadway.Message
def start_link(_opts) do
Broadway.start_link(__MODULE__,
name: __MODULE__,
producer: [
module: {BroadwaySQS.Producer,
queue_url: System.fetch_env!("SQS_QUEUE_URL"),
config: [access_key_id: System.fetch_env!("AWS_ACCESS_KEY_ID"),
secret_access_key: System.fetch_env!("AWS_SECRET_ACCESS_KEY"),
region: "ap-northeast-1"]},
concurrency: 1
],
processors: [default: [concurrency: 10, max_demand: 5]],
batchers: [default: [concurrency: 4, batch_size: 100, batch_timeout: 1_000]]
)
end
@impl true
def handle_message(_, %Message{data: raw} = msg, _ctx) do
case Jason.decode(raw) do
{:ok, payload} -> Message.put_data(msg, payload)
{:error, _} -> Message.failed(msg, :invalid_json)
end
end
@impl true
def handle_batch(:default, messages, _batch_info, _ctx) do
MyApp.Repo.insert_all("events", Enum.map(messages, & &1.data), on_conflict: :nothing)
messages
end
end
5.3 与 Kafka/RabbitMQ 对接
| 数据源 | 生产者模块 | 关键配置 |
|---|---|---|
| SQS | BroadwaySQS.Producer | queue_url、wait_time_seconds、visibility_timeout |
| Kafka | BroadwayKafka.Producer | hosts、group_id、topics、partition 分配策略 |
| RabbitMQ | BroadwayRabbitMQ.Producer | queue、declare、on_success、qos |
| Redis Streams | BroadwayCloudPubSub / 第三方 | stream、group、consumer |
Kafka 版本的处理器需要感知 partition 以便按分区顺序处理,典型配置是 processors: [default: [concurrency: 10, max_demand: 10]] 配合 batchers: [default: [concurrency: 1, batch_size: 500, batch_timeout: 2_000]]——批处理器并发设为 1,保证同一分区内的批次不会乱序。
5.4 消息生命周期
Broadway 的每条消息都带一个「确认器(acknowledger)」,成功处理才向源端确认。Message.put_data/2 默认在消息成功返回后自动 ack;若需要异步确认(比如先落库再等外部回调),可用 Message.put_acknowledger(msg, fn ack, _ -> ack.() end) 手动控制时机;失败则用 Message.failed(msg, {:http_error, 500}) 交给 handle_failed/2。
这一步是「至少一次」语义的基础:只要消息没被 ack,源端就会重新投递。因此 handle_batch/4 里的写库操作必须幂等——用唯一索引 + on_conflict: :nothing,或先写幂等键再执行。
六、错误处理、重启与可观测
6.1 失败重试与死信
@impl true
def handle_failed(messages, _ctx) do
Enum.map(messages, fn msg ->
if msg.metadata.retry_count < 3 do
Message.retry(msg) # 退避重投,要求幂等
else
MyApp.DeadLetter.record(msg) # 超阈值落死信,等待人工介入
Message.failed(msg, :exhausted)
end
end)
end
| 失败类型 | 处理策略 | 注意点 |
|---|---|---|
| 瞬时错误(网络抖动、锁超时) | Message.retry/1 退避重试 | 必须幂等 |
| 数据错误(JSON 解析失败) | 立即 Message.failed/2 | 重试无意义,直接死信 |
| 下游过载(DB 连接池耗尽) | 重试 + 降低 max_demand | 让背压生效,而非盲目重试 |
| 毒丸消息 | 死信队列 + 告警 | 否则会无限循环阻塞分区 |
6.2 与监督树集成
Broadway 管道本身就是一个监督树,启动时应挂在应用的 supervisor 下:
children = [
MyApp.Repo,
{MyApp.SQSPipeline, []}
]
Supervisor.start_link(children, strategy: :one_for_one, name: MyApp.Supervisor)
processors 与 batchers 各自是独立的 GenStage 进程,任何一个崩溃都会被 Broadway 内部的 supervisor 重启,管道其余部分不受影响——这正是 OTP 监督树「故障隔离」思想在数据管道上的体现(参见 https://plumephp.com/elixir-otp-supervision-tasks/)。
6.3 Telemetry 事件
Broadway 内置了完整的 Telemetry 事件,可以零侵入接入指标系统:
| 事件 | 触发时机 | 关键测量值 |
|---|---|---|
[:broadway, :processor, :message, :start] | 处理一条消息开始 | system_time |
[:broadway, :processor, :message, :stop] | 处理成功 | duration |
[:broadway, :processor, :message, :exception] | 处理抛异常 | kind、reason |
[:broadway, :batcher, :batch, :stop] | 一批处理完成 | batch_size、duration |
:telemetry.attach_many("broadway-metrics",
[[:broadway, :processor, :message, :stop],
[:broadway, :processor, :message, :exception]],
fn event, measurements, metadata, _cfg ->
MyApp.Metrics.record(event, measurements, metadata)
end, nil)
6.4 关键指标与告警
生产环境至少应监控以下几项,任何一项异常都意味着管道正在劣化:
- 消息积压量:源端
ApproximateNumberOfMessages(SQS)或 consumer lag(Kafka); - 处理延迟 P99:
[:broadway, :processor, :message, :stop]的duration; - 失败率:
exception事件数除以stop事件数; - 批大小分布:长期低于
batch_size说明上游流量不足或batch_timeout太短; - 进程内存:processor 进程堆增长往往意味着单条消息过大或存在状态泄漏。
七、最佳实践与总结
- 需求参数按「处理代价」调优:
max_demand应约等于「单条处理耗时 × 目标批延迟」;重活小批次、轻活大批次,避免用一套配置打天下; - 聚合前必须分区:Flow 的
reduce依赖分区,漏掉Flow.partition/2会得到错误的并行结果; - 常驻管道用 Broadway,一次性作业用 Flow:前者有 ack/重试/死信,后者胜在语法简洁;
- 幂等是重试的前提:所有
handle_message中的副作用(写库、发消息、调外部 API)都要可重复执行,配合唯一约束或幂等键; - 背压要靠消费端:不要在生产者端 sleep 或丢弃,让需求自然收敛,才是 BEAM 生态的「正确姿势」;
- 可观测先行:上线前接好 Telemetry,把积压量、延迟、失败率画成看板,见 https://plumephp.com/erlang-logging-telemetry-observability/ 的指标章节;
- 消息中间件选型:队列语义(ack、死信、顺序)差异很大,RabbitMQ 与 Kafka 的取舍见 https://plumephp.com/erlang-rabbitmq-messaging/。
流式处理的本质是一场关于「速率」的博弈。GenStage 用需求驱动的拉模型把这场博弈变成了架构的自发属性,Flow 让并行计算回到函数式管道,Broadway 则把生产必需的 ack、重试、死信、可观测一并打包。掌握这三层,就能在 BEAM 上构建出「流量翻十倍也不崩」的数据管道——它不需要复杂的限流配置,因为它从设计上就慢不下来。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。