能否Mock Kafka环境本地运行Flink JUnit测试?附Job测试与依赖问题
Flink JUnit测试:Mock Kafka、Job测试及版本兼容解决方案
一、能不能Mock Kafka做本地Flink JUnit测试?
当然可以,常用两种方案:
- 嵌入式Kafka:借助
kafka-streams-test-utils里的EmbeddedKafkaBroker,在测试时启动一个本地临时Kafka集群 - Flink测试工具类:用
SourceTestHarness和SinkTestHarness直接Mock Kafka源和Sink,无需启动真实Kafka服务
二、你的Flink Job测试方案
先梳理你Job的核心逻辑:从Kafka消费数据→滑动窗口处理→通过Socket调用外部服务→输出结果到Kafka。要测试这个Job,得先重构代码降低耦合,再Mock依赖服务验证流程。
第一步:重构代码提升可测试性
原代码硬编码了Kafka、Socket地址,且在RichAllWindowFunction内直接创建Socket,完全无法Mock。需做以下改动:
- 将Kafka、Socket配置改为可注入参数
- 把Socket通信逻辑抽成独立服务类,方便测试时替换为Mock实现
重构后的代码示例:
@Component public class RedisController { @Value("${flink.consumer.topic.name}") private String flinkConsumerTopic; @Value("${flink.kafka.bootstrap.servers:localhost:9092}") private String kafkaBootstrapServers; @Value("${socket.server.address:localhost}") private String socketServerAddress; @Value("${socket.server.port:8888}") private int socketServerPort; private final SocketService socketService; // 构造注入Socket服务,测试时可传入Mock实现 public RedisController(SocketService socketService) { this.socketService = socketService; } public void run() throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); Properties properties = new Properties(); properties.setProperty("bootstrap.servers", kafkaBootstrapServers); FlinkKafkaConsumer<ObjectNode> consumer = new FlinkKafkaConsumer<>(flinkConsumerTopic, new JSONKeyValueDeserializationSchema(false), properties); DataStream<ObjectNode> stream = env.addSource(consumer); DataStream<String> windowedStream = stream .windowAll(SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(1))) .apply(new RichAllWindowFunction<ObjectNode, String, TimeWindow>() { @Override public void open(Configuration parameters) throws Exception { socketService.connect(socketServerAddress, socketServerPort); } @Override public void apply(TimeWindow timeWindow, Iterable<ObjectNode> input, Collector<String> out) throws Exception { String inputJson = input.toString(); String response = socketService.sendAndReceive(inputJson); out.collect(response); } @Override public void close() throws Exception { socketService.disconnect(); } }); FlinkKafkaProducer<String> kafkaProducer = new FlinkKafkaProducer<>( "output_topic", new SimpleStringSchema(), properties); windowedStream.addSink(kafkaProducer); env.execute("Sliding Window Kafka Job"); } } // 抽象Socket服务接口 public interface SocketService { void connect(String address, int port) throws IOException; String sendAndReceive(String input) throws IOException; void disconnect() throws IOException; } // 生产环境实现类 public class DefaultSocketService implements SocketService { private Socket socket; private OutputStream outputStream; private InputStream inputStream; @Override public void connect(String address, int port) throws IOException { socket = new Socket(address, port); outputStream = socket.getOutputStream(); inputStream = socket.getInputStream(); } @Override public String sendAndReceive(String input) throws IOException { byte[] inputBytes = input.getBytes(); outputStream.write(inputBytes); outputStream.flush(); byte[] responseBytes = new byte[15034]; int bytesRead = inputStream.read(responseBytes); return bytesRead != -1 ? new String(responseBytes, 0, bytesRead) : ""; } @Override public void disconnect() throws IOException { if (inputStream != null) inputStream.close(); if (outputStream != null) outputStream.close(); if (socket != null) socket.close(); } }
第二步:编写JUnit测试用例
分两种测试场景:单算子测试(聚焦窗口逻辑)和端到端测试(模拟完整流程)。
场景1:单算子测试(用Flink TestHarness)
Mock Socket服务,仅测试窗口处理逻辑,无需启动Kafka:
import org.apache.flink.streaming.api.windowing.time.Time; import org.apache.flink.streaming.api.windowing.windows.TimeWindow; import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.mockito.Mock; import org.mockito.MockitoAnnotations; import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.JsonNodeFactory; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.when; import static org.junit.jupiter.api.Assertions.*; public class RedisControllerWindowTest { @Mock private SocketService mockSocketService; private OneInputStreamOperatorTestHarness<ObjectNode, String> testHarness; @BeforeEach void setUp() throws Exception { MockitoAnnotations.openMocks(this); // 初始化窗口算子 RichAllWindowFunction<ObjectNode, String, TimeWindow> windowFunction = new RichAllWindowFunction<>() { @Override public void open(org.apache.flink.configuration.Configuration parameters) throws Exception { mockSocketService.connect("localhost", 8888); } @Override public void apply(TimeWindow timeWindow, Iterable<ObjectNode> input, Collector<String> out) throws Exception { String response = mockSocketService.sendAndReceive(input.toString()); out.collect(response); } }; // 创建测试Harness,模拟窗口算子运行环境 testHarness = new OneInputStreamOperatorTestHarness<>( new WindowAllOperator<>( SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(1)), new InternalTimeWindow.TimestampedValueSerializer<>(new org.apache.flink.api.common.serialization.Serializer<ObjectNode>() { @Override public void serialize(ObjectNode element, OutputStream outputStream) throws IOException { new com.fasterxml.jackson.databind.ObjectMapper().writeValue(outputStream, element); } @Override public ObjectNode deserialize(InputStream inputStream) throws IOException { return new com.fasterxml.jackson.databind.ObjectMapper().readValue(inputStream, ObjectNode.class); } }), windowFunction ) ); testHarness.open(); } @AfterEach void tearDown() throws Exception { testHarness.close(); } @Test void testWindowProcessing() throws Exception { // Mock Socket服务的返回值 when(mockSocketService.sendAndReceive(anyString())).thenReturn("mock_response"); // 构造测试数据 ObjectNode testNode = JsonNodeFactory.instance.objectNode(); testNode.put("patientId", "123"); // 发送数据到算子,指定事件时间 testHarness.processElement(testNode, 1000L); // 推进处理时间到窗口触发(窗口大小10秒,推进到11000毫秒) testHarness.setProcessingTime(11000L); // 验证输出结果 var outputValues = testHarness.extractOutputValues(); assertEquals(1, outputValues.size()); assertEquals("mock_response", outputValues.get(0)); } }
场景2:端到端测试(用嵌入式Kafka)
启动临时Kafka集群,模拟完整的生产→消费→处理→输出流程:
import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.TestInstance; import org.springframework.kafka.test.EmbeddedKafkaBroker; import org.springframework.kafka.test.context.EmbeddedKafka; import org.springframework.kafka.test.utils.KafkaTestUtils; import java.io.BufferedReader; import java.io.IOException; import java.io.InputStreamReader; import java.io.PrintWriter; import java.net.ServerSocket; import java.net.Socket; import static org.junit.jupiter.api.Assertions.assertEquals; @TestInstance(TestInstance.Lifecycle.PER_CLASS) @EmbeddedKafka(partitions = 1, topics = {"test_input_topic", "output_topic"}, ports = {9092}) public class RedisControllerE2ETest { private EmbeddedKafkaBroker embeddedKafka; @BeforeAll void setUp(EmbeddedKafkaBroker broker) { this.embeddedKafka = broker; } @AfterAll void tearDown() { embeddedKafka.destroy(); } @Test void testEndToEndJob() throws Exception { // 启动本地Socket服务器,模拟外部服务响应 new Thread(() -> { try (ServerSocket serverSocket = new ServerSocket(8888)) { Socket clientSocket = serverSocket.accept(); BufferedReader in = new BufferedReader(new InputStreamReader(clientSocket.getInputStream())); PrintWriter out = new PrintWriter(clientSocket.getOutputStream(), true); // 读取Flink发送的请求,返回固定响应 String input = in.readLine(); out.println("e2e_test_response"); } catch (IOException e) { // 忽略测试结束时的Socket关闭异常 } }).start(); // 创建Job实例,替换配置为嵌入式Kafka地址 RedisController controller = new RedisController(new DefaultSocketService()); controller.flinkConsumerTopic = "test_input_topic"; controller.kafkaBootstrapServers = embeddedKafka.getBrokersAsString(); // 异步启动Job,避免阻塞测试线程 new Thread(() -> { try { controller.run(); } catch (Exception e) { e.printStackTrace(); } }).start(); // 生产测试数据到输入Topic KafkaTestUtils.sendMessages(embeddedKafka, "test_input_topic", "{\"patientId\":\"456\"}"); // 从输出Topic消费结果并验证 String output = KafkaTestUtils.getSingleRecord(embeddedKafka, "output_topic").value(); assertEquals("e2e_test_response", output); } }
三、依赖版本兼容问题解决
版本冲突是Flink测试的常见坑,按以下步骤解决:
- 对齐Flink与Kafka版本:Flink对Kafka客户端版本有严格要求,参考官方兼容表:
- Flink 1.16.x → Kafka 2.8.x ~ 3.2.x
- Flink 1.17.x → Kafka 3.0.x ~ 3.4.x
- Flink 1.18.x → Kafka 3.2.x ~ 3.5.x
- 排除冲突依赖:如果项目中其他依赖引入了不兼容的Kafka客户端版本,用Maven的
<exclusions>排除:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>${flink.version}</version> <exclusions> <exclusion> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> </exclusion> </exclusions> </dependency> <!-- 手动添加兼容版本的kafka-clients --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>${kafka.version}</version> </dependency>
- 用BOM统一版本管理:在Maven的
<dependencyManagement>里引入Flink BOM,自动管理所有Flink相关依赖的版本:
<dependencyManagement> <dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-bom</artifactId> <version>${flink.version}</version> <type>pom</type> <scope>import</scope> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>${kafka.version}</version> </dependency> </dependencies> </dependencyManagement>
内容的提问来源于stack exchange,提问作者Anand SL
相关产品推荐
相关产品推荐

