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

如何测试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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 13:20:51