You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何测试Elixir GenStage Consumer组件?

我明白你现在的困扰——测试Producer的时候用dummy consumer顺风顺水,但轮到测试Consumer就卡壳了。别担心,咱们一步步拆解问题,给你一套实用的测试方案。

首先,咱们先理清楚你的DataConsumer核心逻辑:它订阅了DataProducer,只处理50 < n < 100的事件,并且每次最多请求10个事件。测试的关键就是要验证这两个行为:是否正确筛选事件,是否遵守demand规则。

第一步:改造Consumer,让它更易测试

原来的DataConsumer硬编码订阅了DataProducer,测试时很难替换依赖。咱们先把它改成可配置的:

defmodule DataConsumer do
  use GenStage

  # 新增参数:允许传入要订阅的Producer,以及测试进程的PID(用于接收处理结果)
  def start_link(producer, test_owner \\ nil) do
    GenStage.start_link(__MODULE__, {producer, test_owner})
  end

  def init({producer, test_owner}) do
    {:consumer, test_owner, 
     subscribe_to: [{producer, selector: fn n -> n > 50 && n < 100 end, max_demand: 10}]}
  end

  def handle_events(events, _from, test_owner) do
    filtered_events = Enum.filter(events, fn n -> n > 50 && n < 100 end)
    
    # 保留原有日志逻辑
    for event <- filtered_events do
      Logger.info inspect({self(), event, test_owner})
    end

    # 如果是测试场景,把处理后的事件发送回测试进程
    if test_owner do
      send(test_owner, {:consumer_processed, filtered_events})
    end

    {:noreply, [], test_owner}
  end
end

第二步:创建测试用Producer

我们需要一个能手动控制事件输出的Producer,代替真实的DataProducer,这样才能精确测试Consumer的筛选逻辑:

defmodule TestProducer do
  use GenStage

  def start_link(initial_events \\ []) do
    GenStage.start_link(__MODULE__, {:queue.from_list(initial_events), 0}, name: __MODULE__)
  end

  # 外部接口:添加事件到Producer队列
  def add_events(stage, events) do
    GenStage.call(stage, {:add_events, events})
  end

  def init({queue, sent_count}) do
    {:producer, {queue, sent_count}}
  end

  def handle_call({:add_events, events}, _from, {queue, sent}) do
    new_queue = Enum.reduce(events, queue, &Queue.in/2)
    {:reply, :ok, [], {new_queue, sent}}
  end

  def handle_demand(demand, {queue, sent}) do
    {events_to_send, new_queue} = Queue.split(queue, demand)
    {:noreply, events_to_send, {new_queue, sent + length(events_to_send)}}
  end
end

第三步:编写Consumer测试用例

现在我们可以编写两种测试场景:验证筛选逻辑,以及验证demand规则。

测试1:验证Consumer只处理符合条件的事件

defmodule DataConsumerTest do
  use ExUnit.Case

  test "only processes events between 50 and 100" do
    # 1. 启动测试用Producer
    {:ok, test_producer} = TestProducer.start_link()
    # 2. 启动Consumer,传入测试进程PID(用于接收结果)
    test_owner = self()
    {:ok, _consumer} = DataConsumer.start_link(test_producer, test_owner)

    # 3. 给Producer添加混合事件:符合条件的+不符合的
    TestProducer.add_events(test_producer, [45, 60, 90, 105, 75])
    # 触发Producer发送事件
    GenStage.demand(test_producer, 5)

    # 4. 断言Consumer返回的结果是筛选后的
    assert_receive {:consumer_processed, [60, 90, 75]}
    # 确认不符合条件的事件被忽略
    refute_receive {:consumer_processed, [45, 105]}
  end

测试2:验证Consumer的max_demand规则

test "respects max_demand setting" do
    {:ok, test_producer} = TestProducer.start_link()
    test_owner = self()
    {:ok, _consumer} = DataConsumer.start_link(test_producer, test_owner)

    # 添加超过max_demand(10)的事件
    TestProducer.add_events(test_producer, Enum.to_list(1..15))
    GenStage.demand(test_producer, 15)

    # 第一次只会收到10个符合条件的事件(51-60)
    assert_receive {:consumer_processed, Enum.to_list(51..60)}
    # 继续触发demand,才会收到剩下的符合条件的事件(61-65)
    GenStage.demand(test_producer, 5)
    assert_receive {:consumer_processed, Enum.to_list(61..65)}
  end
end

核心思路总结

  • 控制输入:用测试用Producer代替真实依赖,手动控制发送的事件,确保测试场景可复现。
  • 验证输出:让Consumer把处理结果发送回测试进程,或者通过ExUnit.CaptureLog捕获日志,直接验证行为是否符合预期。
  • 依赖解耦:让Consumer的订阅目标可配置,避免测试时依赖真实的DataProducer,让测试更独立。

内容的提问来源于stack exchange,提问作者Mr H

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 08:52:29