如何测试Flink的RichFilterFunction?含状态算子测试需求
测试Flink RichFilterFunction的方案
我想找到测试RichFilterFunction的方法,就像测试RichMapFunction那样,但全网搜索没找到解决方案。以下是示例代码:
public class Person{ public String name; public Integer age; } public class OCRKafkaFilter extends RichFilterFunction<Person> { @Override public boolean filter(Person person) { return person.age > 21; } }
请问如何为这类带状态的富过滤函数编写测试用例?
核心思路
测试RichFilterFunction和测试其他Rich算子逻辑一致,利用Flink官方提供的OneInputStreamOperatorTestHarness(无键控场景)或KeyedOneInputStreamOperatorTestHarness(键控场景)模拟Flink运行环境,触发算子的生命周期方法(open/close等),并验证输入输出及状态变化。
步骤1:引入测试依赖
确保项目包含Flink测试工具依赖(以Maven为例):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-test-utils</artifactId> <version>${flink.version}</version> <scope>test</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>${flink.version}</version> <scope>test</scope> <type>test-jar</type> </dependency>
步骤2:编写带状态的RichFilterFunction示例
先扩展示例,增加状态逻辑(比如统计被过滤的人数),以便体现状态测试的核心:
public class StatefulOCRKafkaFilter extends RichFilterFunction<Person> { private ValueState<Integer> filteredCountState; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 注册状态 ValueStateDescriptor<Integer> descriptor = new ValueStateDescriptor<>( "filteredCount", Integer.class, 0 // 默认值 ); filteredCountState = getRuntimeContext().getState(descriptor); } @Override public boolean filter(Person person) throws Exception { boolean shouldKeep = person.age > 21; if (!shouldKeep) { // 更新状态:过滤计数+1 filteredCountState.update(filteredCountState.value() + 1); } return shouldKeep; } // 供测试获取状态用(可选,也可通过TestHarness直接获取) public ValueState<Integer> getFilteredCountState() { return filteredCountState; } }
步骤3:编写测试用例
使用OneInputStreamOperatorTestHarness完成带状态的Filter测试:
import org.apache.flink.api.common.state.ValueState; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.operators.SimpleOperatorFactory; import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness; import org.junit.Before; import org.junit.Test; import static org.junit.Assert.*; public class StatefulOCRKafkaFilterTest { private OneInputStreamOperatorTestHarness<Person, Person> testHarness; private StatefulOCRKafkaFilter filterFunction; @Before public void setup() throws Exception { // 实例化Filter函数 filterFunction = new StatefulOCRKafkaFilter(); // 创建算子测试容器 testHarness = new OneInputStreamOperatorTestHarness<>( SimpleOperatorFactory.of(filterFunction).createStreamOperator(new Configuration()) ); // 初始化测试容器,触发算子的open方法 testHarness.open(); } @Test public void testFilterAndStateUpdate() throws Exception { // 输入测试数据:3个Person,其中2个年龄<=21会被过滤 testHarness.processElement(new Person("Alice", 20)); testHarness.processElement(new Person("Bob", 25)); testHarness.processElement(new Person("Charlie", 18)); // 验证输出:只有Bob被保留 assertEquals(1, testHarness.getOutput().size()); assertEquals("Bob", testHarness.getOutput().poll().name); // 验证状态:过滤计数应为2 assertEquals(Integer.valueOf(2), filterFunction.getFilteredCountState().value()); // 或者通过TestHarness直接获取状态(无需修改原函数) ValueState<Integer> stateFromHarness = testHarness.getStateBackend().getKeyedStateStore() .getState(new ValueStateDescriptor<>("filteredCount", Integer.class)); assertEquals(Integer.valueOf(2), stateFromHarness.value()); } @Test public void testFilterKeepPerson() throws Exception { testHarness.processElement(new Person("David", 30)); assertEquals(1, testHarness.getOutput().size()); assertEquals("David", testHarness.getOutput().poll().name); // 状态计数应为0 assertEquals(Integer.valueOf(0), filterFunction.getFilteredCountState().value()); } }
关键说明
OneInputStreamOperatorTestHarness会模拟Flink运行环境,自动处理算子的生命周期(open/close等),无需手动调用RichFunction的open方法。- 如果是键控场景(Filter前做了keyBy),则使用
KeyedOneInputStreamOperatorTestHarness,需要指定键选择器和类型信息。 - 测试无状态的RichFilterFunction时,只需验证输入输出即可,区别在于TestHarness会触发生命周期方法,确保RichFunction的初始化逻辑被执行。
内容的提问来源于stack exchange,提问作者Praveen Kumar Shah
相关产品推荐
相关产品推荐

