Erlang/Elixir 流式数据处理:GenStage、Flow 与 Broadway 实战

从 GenStage 的需求驱动背压模型出发,系统讲解生产者/消费者/生产者消费者三类角色的回调契约、Flow 的分区窗口与聚合算子,以及 Broadway 构建 SQS/Kafka/RabbitMQ 生产管道的并发编排、错误处理与可观测方案。

在 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_demand1000一次最多向上游索取的事件数
min_demanddiv(max_demand, 2)缓冲区低于此值时触发补货
buffer_sizemax_demand 的一半缓冲上限,超出会丢弃
dispatcherGenStage.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 作业FlowFlow.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 对接

数据源生产者模块关键配置
SQSBroadwaySQS.Producerqueue_url、wait_time_seconds、visibility_timeout
KafkaBroadwayKafka.Producerhosts、group_id、topics、partition 分配策略
RabbitMQBroadwayRabbitMQ.Producerqueue、declare、on_success、qos
Redis StreamsBroadwayCloudPubSub / 第三方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 上构建出「流量翻十倍也不崩」的数据管道——它不需要复杂的限流配置,因为它从设计上就慢不下来。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「erlang」更多文章

  1. Erlang/Elixir 容器化与集群部署:Release、Docker 与 libcluster
  2. Elixir 认证授权实战:JWT、Guardian 与 Phoenix.Token
  3. Elixir HTTP 客户端与连接池:Mint、Finch 与 Req 实战