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

能否Mock Kafka环境本地运行Flink JUnit测试?附Job测试与依赖问题

当然可以,常用两种方案:

  • 嵌入式Kafka:借助kafka-streams-test-utils里的EmbeddedKafkaBroker,在测试时启动一个本地临时Kafka集群
  • Flink测试工具类:用SourceTestHarness和SinkTestHarness直接Mock Kafka源和Sink,无需启动真实Kafka服务

先梳理你Job的核心逻辑:从Kafka消费数据→滑动窗口处理→通过Socket调用外部服务→输出结果到Kafka。要测试这个Job,得先重构代码降低耦合,再Mock依赖服务验证流程。

第一步:重构代码提升可测试性

原代码硬编码了Kafka、Socket地址,且在RichAllWindowFunction内直接创建Socket,完全无法Mock。需做以下改动:

  1. 将Kafka、Socket配置改为可注入参数
  2. 把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测试的常见坑,按以下步骤解决:

  1. 对齐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
  2. 排除冲突依赖:如果项目中其他依赖引入了不兼容的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>
  1. 用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 16:22:37