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

如何对基于FlinkKafkaConsumer011与FlinkKafkaProducer011的Kafka-Flink集成流程做单元/集成测试?

Great question! Testing Kafka-Flink integrations can feel a bit daunting at first, but splitting it into unit tests (no real Kafka needed) and integration tests (with a controlled Kafka instance) gives you confidence in both your business logic and the end-to-end pipeline. Let's break this down for your uppercase transformation flow.


The goal here is to test your core transformation logic without relying on a running Kafka cluster. We'll use Flink's built-in test utilities to mock sources and sinks, focusing on whether your processing works as expected in a Flink context.

Step 1: Add Required Dependencies

First, make sure you have Flink's test utilities in your build (Maven example):

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-test-utils</artifactId>
    <version>${flink.version}</version>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.junit.jupiter</groupId>
    <artifactId>junit-jupiter-api</artifactId>
    <version>${jupiter.version}</version>
    <scope>test</scope>
</dependency>

Step 2: Write the Unit Test

We'll use a MiniClusterResource to simulate a Flink cluster, a CollectionSource to feed test input strings, and a custom TestSink to capture output. Then we'll verify if the output is correctly uppercased.

Here's a code example:

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.test.util.MiniClusterResource;
import org.junit.ClassRule;
import org.junit.Test;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import static org.junit.Assert.assertEquals;

public class UppercaseTransformationTest {

    @ClassRule
    public static MiniClusterResource flinkCluster = new MiniClusterResource.Builder()
            .setNumberSlotsPerTaskManager(1)
            .setNumberTaskManagers(1)
            .build();

    @Test
    public void testUppercaseTransformation() throws Exception {
        // 1. Set up Flink environment
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1); // Keep parallelism low for deterministic results

        // 2. Create test input
        List<String> inputStrings = Arrays.asList("hello", "kafka-flink", "test");
        DataStream<String> inputStream = env.fromCollection(inputStrings);

        // 3. Apply your transformation (reuse the same logic as production)
        DataStream<String> uppercaseStream = inputStream.map(String::toUpperCase);

        // 4. Capture output with a custom TestSink
        List<String> output = new ArrayList<>();
        uppercaseStream.addSink(new TestSink<>(output));

        // 5. Execute the test job
        env.execute("Uppercase Test Job");

        // 6. Verify results match expectations
        List<String> expectedOutput = Arrays.asList("HELLO", "KAFKA-FLINK", "TEST");
        assertEquals(expectedOutput, output);
    }

    // Simple sink to collect test output
    private static class TestSink<T> implements SinkFunction<T> {
        private final List<T> outputCollection;

        public TestSink(List<T> outputCollection) {
            this.outputCollection = outputCollection;
        }

        @Override
        public void invoke(T value, Context context) {
            outputCollection.add(value);
        }
    }
}

This test isolates your core logic—you're not touching Kafka at all, just validating that the uppercase mapping works correctly within a Flink pipeline.


Integration Testing (End-to-End with Kafka)

Now let's test the full pipeline: Kafka → Flink → Kafka. We'll use an embedded Kafka instance (either via Kafka's test modules or Testcontainers) to avoid relying on external infrastructure.

Option 1: Using Kafka's Embedded Cluster

Kafka provides an embedded cluster for testing via its test artifact.

Dependencies

Add this to your build:

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka_2.13</artifactId>
    <version>${kafka.version}</version>
    <classifier>test</classifier>
    <scope>test</scope>
</dependency>

Test Code Example

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer011;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer011;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.apache.kafka.common.serialization.StringSerializer;
import org.apache.kafka.streams.integration.utils.EmbeddedKafkaCluster;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.Properties;
import java.util.concurrent.TimeUnit;
import static org.junit.Assert.assertEquals;

public class KafkaFlinkIntegrationTest {

    private static final String INPUT_TOPIC = "input-topic";
    private static final String OUTPUT_TOPIC = "output-topic";
    private EmbeddedKafkaCluster kafkaCluster;

    @Before
    public void startKafka() throws Exception {
        // Configure embedded Kafka
        Properties kafkaProps = new Properties();
        kafkaProps.setProperty("broker.id", "1");
        kafkaProps.setProperty("listeners", "PLAINTEXT://localhost:9092");
        kafkaProps.setProperty("num.partitions", "1");

        kafkaCluster = new EmbeddedKafkaCluster(1, kafkaProps);
        kafkaCluster.start();
        // Create test topics
        kafkaCluster.createTopic(INPUT_TOPIC);
        kafkaCluster.createTopic(OUTPUT_TOPIC);
    }

    @After
    public void stopKafka() {
        if (kafkaCluster != null) {
            kafkaCluster.stop();
        }
    }

    @Test
    public void testEndToEndPipeline() throws Exception {
        // 1. Produce test data to input topic
        Properties producerProps = new Properties();
        producerProps.setProperty("bootstrap.servers", "localhost:9092");
        producerProps.setProperty("key.serializer", StringSerializer.class.getName());
        producerProps.setProperty("value.serializer", StringSerializer.class.getName());

        try (KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps)) {
            producer.send(new ProducerRecord<>(INPUT_TOPIC, "hello world")).get(5, TimeUnit.SECONDS);
            producer.send(new ProducerRecord<>(INPUT_TOPIC, "flink kafka")).get(5, TimeUnit.SECONDS);
        }

        // 2. Set up Flink pipeline with real Kafka connectors
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // Configure Kafka Consumer
        Properties consumerProps = new Properties();
        consumerProps.setProperty("bootstrap.servers", "localhost:9092");
        consumerProps.setProperty("group.id", "test-group");
        consumerProps.setProperty("auto.offset.reset", "earliest");
        consumerProps.setProperty("key.deserializer", StringDeserializer.class.getName());
        consumerProps.setProperty("value.deserializer", StringDeserializer.class.getName());

        FlinkKafkaConsumer011<String> consumer = new FlinkKafkaConsumer011<>(
                INPUT_TOPIC,
                new SimpleStringSchema(),
                consumerProps
        );

        DataStream<String> inputStream = env.addSource(consumer);
        DataStream<String> uppercaseStream = inputStream.map(String::toUpperCase);

        // Configure Kafka Producer
        Properties producerPropsFlink = new Properties();
        producerPropsFlink.setProperty("bootstrap.servers", "localhost:9092");
        producerPropsFlink.setProperty("key.serializer", StringSerializer.class.getName());
        producerPropsFlink.setProperty("value.serializer", StringSerializer.class.getName());

        uppercaseStream.addSink(new FlinkKafkaProducer011<>(
                OUTPUT_TOPIC,
                new SimpleStringSchema(),
                producerPropsFlink
        ));

        // 3. Run Flink job in a separate thread (so we can consume output)
        new Thread(() -> {
            try {
                env.execute("Kafka-Flink Uppercase Job");
            } catch (Exception e) {
                e.printStackTrace();
            }
        }).start();

        // 4. Consume from output topic and verify results
        Properties consumerOutputProps = new Properties();
        consumerOutputProps.setProperty("bootstrap.servers", "localhost:9092");
        consumerOutputProps.setProperty("group.id", "test-output-group");
        consumerOutputProps.setProperty("auto.offset.reset", "earliest");
        consumerOutputProps.setProperty("key.deserializer", StringDeserializer.class.getName());
        consumerOutputProps.setProperty("value.deserializer", StringDeserializer.class.getName());

        try (KafkaConsumer<String, String> outputConsumer = new KafkaConsumer<>(consumerOutputProps)) {
            outputConsumer.subscribe(Collections.singletonList(OUTPUT_TOPIC));
            // Wait for messages to be processed
            ConsumerRecords<String, String> records = outputConsumer.poll(Duration.ofSeconds(10));
            
            List<String> receivedMessages = new ArrayList<>();
            records.forEach(record -> receivedMessages.add(record.value()));

            List<String> expected = Arrays.asList("HELLO WORLD", "FLINK KAFKA");
            assertEquals(expected, receivedMessages);
        }
    }
}

Option 2: Using Testcontainers (More Realistic Testing)

If you want a closer-to-production Kafka instance, Testcontainers lets you spin up a Dockerized Kafka cluster in your tests. Add the dependency first:

<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>kafka</artifactId>
    <version>${testcontainers.version}</version>
    <scope>test</scope>
</dependency>

The test logic will mirror the embedded cluster example—you'll just use KafkaContainer to start/stop the cluster instead of EmbeddedKafkaCluster.


Key Testing Tips

  • Keep Tests Deterministic: Use parallelism=1 in tests to avoid message ordering issues.
  • Clean Up Resources: Always stop clusters/containers after tests to prevent resource leaks.
  • Use auto.offset.reset=earliest: Ensures your test consumer picks up all messages produced before it starts.
  • Isolate Tests: Create unique topics per test to avoid cross-test contamination.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:11:09