如何对基于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.
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

