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

单元测试含ProcessContext的DoFn:如何实例化ProcessContext对象?

如何实例化DoFn.ProcessContext进行ParDo处理器测试?

你遇到的错误是因为DoFn.ProcessContext是DoFn的非静态内部类,且它的实例由Beam运行时环境负责创建,包含了元素、输出通道、窗口等运行时上下文信息,无法直接通过new关键字手动实例化。下面提供两种可行的测试方案:

方案一:使用Beam官方测试框架(推荐,适合集成测试)

通过TestPipeline和TestStream构建完整的测试管道,模拟输入数据和侧输入,让Beam框架自动管理ProcessContext实例:

import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.testing.TestStream;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.transforms.View;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
import org.junit.Rule;
import org.junit.Test;

public class DoFnHandlerFnTest {
    @Rule
    public final TestPipeline pipeline = TestPipeline.create();

    @Test
    public void testProcessElement() {
        // 1. 构造测试输入元素
        KV<Row, Row> testElement = KV.of(
            Row.withSchema(/* 传入你的Row Schema */).addValues(/* 元素值 */).build(),
            Row.withSchema(/* 传入你的Row Schema */).addValues(/* 元素值 */).build()
        );

        // 2. 构造侧输入数据
        PCollection<Row> si1Source = pipeline.apply(
            TestStream.create(/* 传入Row的Coder */)
                .addElements(/* 侧输入1的测试数据 */)
                .advanceWatermarkToInfinity()
        );
        PCollection<Row> si2Source = pipeline.apply(
            TestStream.create(/* 传入Row的Coder */)
                .addElements(/* 侧输入2的测试数据 */)
                .advanceWatermarkToInfinity()
        );

        // 3. 转换为侧输入视图
        var si1View = si1Source.apply(View.asMap());
        var si2View = si2Source.apply(View.asList());

        // 4. 创建输入PCollection
        PCollection<KV<Row, Row>> input = pipeline.apply(
            TestStream.create(KV.getKVCoder(/* Row的Coder */, /* Row的Coder */))
                .addElements(testElement)
                .advanceWatermarkToInfinity()
        );

        // 5. 应用ParDo并传入侧输入
        PCollection<Row> output = input.apply(
            ParDo.of(new DoFnHandlerFn()).withSideInputs(si1View, si2View)
        );

        // 6. 验证输出结果(可结合Beam的断言工具,比如Assert.that())
        // Assert.that(output, Matchers.containsInAnyOrder(/* 预期输出Row */));
    }
}

方案二:用Mockito模拟ProcessContext(适合单元测试单个方法)

如果只需要测试processElement方法的业务逻辑,可通过Mockito模拟ProcessContext的核心方法(比如element()、output()):

import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.Row;
import org.junit.Test;
import org.mockito.Mockito;

import java.util.List;
import java.util.Map;

public class DoFnHandlerFnTest {
    @Test
    public void testProcessElement() {
        // 1. 构造测试元素
        KV<Row, Row> testElement = KV.of(
            Row.withSchema(/* 传入你的Row Schema */).addValues(/* 元素值 */).build(),
            Row.withSchema(/* 传入你的Row Schema */).addValues(/* 元素值 */).build()
        );

        // 2. 模拟ProcessContext实例
        DoFn<KV<Row, Row>, Row>.ProcessContext mockCtx = Mockito.mock(DoFn.ProcessContext.class);
        Mockito.when(mockCtx.element()).thenReturn(testElement);

        // 3. 构造侧输入测试数据
        Map<Row, Iterable<Row>> si1 = /* 构建测试用的侧输入1数据 */;
        List<Row> si2 = /* 构建测试用的侧输入2数据 */;

        // 4. 调用目标方法
        DoFnHandlerFn doFn = new DoFnHandlerFn();
        doFn.processElement(mockCtx, si1, si2);

        // 5. 验证逻辑(比如检查是否调用了ctx.output(),并传入预期值)
        // Mockito.verify(mockCtx).output(/* 预期输出的Row */);
    }
}

需要注意:如果你的processElement方法中使用了ProcessContext的其他方法(比如output()、timestamp()),需要在Mockito中对应的模拟这些方法的行为。

内容的提问来源于stack exchange,提问作者Shlomi Elbaz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:07:23