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

如何对Flink中带自定义指标的RichSinkFunction进行单元测试?

解决RichSinkFunction中Counter指标的单元测试问题

针对你遇到的Counter无法Mock的问题,可以通过Flink RichFunction自带的setRuntimeContext方法注入Mock依赖,结合Mockito完成单元测试,具体步骤如下:

核心思路

通过Mock RuntimeContext、MetricGroup 和 Counter 构建依赖链,让Sink的open方法初始化时拿到Mock的Counter,从而验证invoke方法中指标的递增逻辑。

测试代码示例

import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.apache.flink.api.common.functions.RuntimeContext;
import org.apache.flink.metrics.Counter;
import org.apache.flink.metrics.MetricGroup;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;

import static org.mockito.Mockito.*;

@ExtendWith(MockitoExtension.class)
class CustomRichSinkFunctionTest {

    @Mock
    private RuntimeContext mockRuntimeContext;

    @Mock
    private MetricGroup mockMetricGroup;

    @Mock
    private Counter mockCounter;

    @Test
    void invoke_shouldIncrementCounter() throws Exception {
        // 初始化自定义Sink实例
        CustomRichSinkFunction sink = new CustomRichSinkFunction();

        // 构建Mock依赖链:RuntimeContext → MetricGroup → Counter
        when(mockRuntimeContext.getMetricGroup()).thenReturn(mockMetricGroup);
        when(mockMetricGroup.counter("myCounter")).thenReturn(mockCounter);

        // 注入Mock的RuntimeContext并触发open方法初始化指标
        sink.setRuntimeContext(mockRuntimeContext);
        sink.open(new Configuration());

        // 构造测试数据并调用invoke方法
        MyData testData = new MyData(); // 根据实际业务构造测试数据
        sink.invoke(testData, mock(SinkFunction.Context.class));

        // 验证Counter的inc方法被调用一次
        verify(mockCounter, times(1)).inc();
    }
}

关键说明

  1. 利用Flink扩展点注入依赖:RichFunction接口提供了setRuntimeContext方法,无需反射即可替换RuntimeContext为Mock实例。
  2. Mock依赖链传递:让Mock的RuntimeContext返回Mock的MetricGroup,再让MetricGroup返回我们控制的Mock Counter,确保open方法初始化的customMetric指向Mock对象。
  3. 验证指标逻辑:通过Mockito的verify方法确认inc()方法的调用次数,直接验证指标递增行为。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 22:03:17