如何对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(); } }
关键说明
- 利用Flink扩展点注入依赖:RichFunction接口提供了
setRuntimeContext方法,无需反射即可替换RuntimeContext为Mock实例。 - Mock依赖链传递:让Mock的RuntimeContext返回Mock的MetricGroup,再让MetricGroup返回我们控制的Mock Counter,确保
open方法初始化的customMetric指向Mock对象。 - 验证指标逻辑:通过Mockito的
verify方法确认inc()方法的调用次数,直接验证指标递增行为。
内容的提问来源于stack exchange,提问作者Olahzzz
相关产品推荐
相关产品推荐

