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

基于scalatest-embedded-kafka的Flink与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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:37:30