基于scalatest-embedded-kafka的Flink与Kafka集成测试咨询
Flink + Kafka 端到端集成测试(基于scalatest-embedded-kafka)
我来帮你完善这个Flink与Kafka的端到端集成测试方案,下面是完整可运行的Scala示例代码,以及关键步骤的说明:
第一步:配置依赖(build.sbt)
确保你的项目引入了必要的依赖:
libraryDependencies ++= Seq( // Flink 核心与流处理依赖(根据你的Flink版本调整) "org.apache.flink" %% "flink-scala" % "1.17.0" % Provided, "org.apache.flink" %% "flink-streaming-scala" % "1.17.0" % Provided, "org.apache.flink" %% "flink-connector-kafka" % "1.17.0", // scalatest-embedded-kafka 依赖 "net.manub" %% "scalatest-embedded-kafka" % "2.8.0" % Test, // ScalaTest 测试框架 "org.scalatest" %% "scalatest" % "3.2.15" % Test )
第二步:完整测试代码
import java.util.Properties import net.manub.embeddedkafka.{EmbeddedKafka, EmbeddedKafkaConfig} import org.apache.flink.api.common.serialization.SimpleStringSchema import org.apache.flink.streaming.api.scala._ import org.apache.flink.streaming.connectors.kafka.{FlinkKafkaConsumer, FlinkKafkaProducer} import org.scalatest.BeforeAndAfterAll import org.scalatest.funsuite.AnyFunSuite class FlinkKafkaE2ETest extends AnyFunSuite with EmbeddedKafka with BeforeAndAfterAll { // 定义测试用的Topic名称 private val inputTopic = "test-input-topic" private val outputTopic = "test-output-topic" // 嵌入式Kafka配置,使用默认端口 private val kafkaConfig = EmbeddedKafkaConfig() override def beforeAll(): Unit = { // 启动嵌入式Kafka,并创建测试用的Topic startEmbeddedKafka(kafkaConfig) createCustomTopic(inputTopic, kafkaConfig) createCustomTopic(outputTopic, kafkaConfig) } override def afterAll(): Unit = { // 停止嵌入式Kafka stopEmbeddedKafka() } test("Flink从Kafka读取数据、处理后写回Kafka的端到端测试") { // 1. 初始化Flink流处理环境 val env = StreamExecutionEnvironment.getExecutionEnvironment // 测试环境设置并行度为1,避免多线程干扰 env.setParallelism(1) // 开启Checkpoint保证Exactly-Once语义(可选,根据你的业务需求) env.enableCheckpointing(1000) // 2. 配置Kafka消费者属性 val consumerProps = new Properties() consumerProps.setProperty("bootstrap.servers", kafkaConfig.kafkaBootstrapServers) consumerProps.setProperty("group.id", "flink-test-group") consumerProps.setProperty("auto.offset.reset", "earliest") // 3. 从Kafka读取数据流 val inputStream: DataStream[String] = env.addSource( new FlinkKafkaConsumer[String](inputTopic, new SimpleStringSchema(), consumerProps) ) // 4. 模拟Flink数据处理逻辑(这里简单做字符串转大写) val processedStream: DataStream[String] = inputStream.map(_.toUpperCase) // 5. 配置Kafka生产者属性 val producerProps = new Properties() producerProps.setProperty("bootstrap.servers", kafkaConfig.kafkaBootstrapServers) // 配置生产者的事务属性(如果开启了Checkpoint,建议开启事务保证Exactly-Once) producerProps.setProperty("transaction.timeout.ms", "60000") // 6. 将处理后的数据写回Kafka processedStream.addSink( new FlinkKafkaProducer[String]( outputTopic, new SimpleStringSchema(), producerProps, FlinkKafkaProducer.Semantic.EXACTLY_ONCE // 保证Exactly-Once语义 ) ) // 7. 发送测试数据到输入Topic val testData = List("hello flink", "kafka integration", "end to end test") publishToKafka(inputTopic, testData)(kafkaConfig) // 8. 异步提交Flink作业 val job = env.executeAsync() // 9. 等待作业处理完成(测试环境下可适当调整等待时间) Thread.sleep(5000) // 取消作业 job.cancel() // 10. 从输出Topic读取数据,验证结果 val consumedData = consumeNumberMessagesFromTopics[Set[String]](Set(outputTopic), testData.size)(kafkaConfig, new SimpleStringSchema()) val expectedData = testData.map(_.toUpperCase).toSet assert(consumedData == expectedData, "处理后的数据与预期不符") } }
关键细节说明
- 嵌入式Kafka生命周期管理:通过继承
EmbeddedKafka特质,结合BeforeAndAfterAll来自动启动/停止Kafka,避免手动管理服务。 - Flink环境配置:测试环境设置并行度为1,减少多线程带来的测试不确定性;开启Checkpoint配合Kafka生产者的
EXACTLY_ONCE语义,保证数据处理的一致性。 - 数据验证:使用
consumeNumberMessagesFromTopics从输出Topic读取指定数量的消息,和预期结果做断言,完成端到端的验证。 - 依赖版本匹配:注意Flink、Kafka客户端、scalatest-embedded-kafka的版本兼容性,避免出现依赖冲突。
内容的提问来源于stack exchange,提问作者Thomas
相关产品推荐
相关产品推荐

