如何测试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
相关产品推荐
相关产品推荐

