单元测试含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
相关产品推荐
相关产品推荐

