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

使用Testcontainers做Kafka+Spark流集成测试如何获取bootstrapServers并创建Topic

问题解决步骤

1 修复依赖配置错误

你现有sbt配置中kafka testcontainers依赖的groupId写错了,将org.dimafeng改为com.dimafeng:

// 错误写法
// "org.dimafeng"     %% "testcontainers-scala-kafka" % "0.39.5" % Test,
// 正确写法
"com.dimafeng"     %% "testcontainers-scala-kafka" % "0.39.5" % Test,

2 获取bootstrapServers

withContainers传入的container默认是泛型容器类型,需要显式转换为KafkaContainer后才能调用bootstrapServers方法:

val kafkaContainer = container.asInstanceOf[KafkaContainer]
val servers = kafkaContainer.bootstrapServers

3 创建指定Kafka Topic

推荐两种常用实现方式:

方式一:在容器内执行官方脚本创建(无需额外依赖)

直接调用容器的execInContainer方法执行kafka内置的topic创建命令:

kafkaContainer.execInContainer(
  "/usr/bin/kafka-topics",
  "--create",
  "--bootstrap-server", "localhost:9092",
  "--replication-factor", "1",
  "--partitions", "1",
  "--topic", "topic1"
)

方式二:使用Kafka AdminClient创建(更灵活,适合批量操作)

先添加kafka-clients测试依赖到sbt:

"org.apache.kafka" % "kafka-clients" % "2.4.1" % Test

再在测试代码中调用AdminClient创建topic:

import org.apache.kafka.clients.admin.{Admin, AdminClientConfig, NewTopic}
import java.util.Properties
import scala.jdk.CollectionConverters._

val props = new Properties()
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, servers)
val admin = Admin.create(props)
admin.createTopics(List(new NewTopic("topic1", 1, 1.toShort)).asJava)
admin.close()

完整修正后的测试代码

import com.dimafeng.testcontainers.{ContainerDef, KafkaContainer}
import com.dimafeng.testcontainers.scalatest.TestContainerForAll
import munit.FunSuite
import org.apache.spark.sql.SparkSession

class Mykafkatest extends FunSuite with TestContainerForAll {
  override val containerDef: ContainerDef = KafkaContainer.Def("confluentinc/cp-kafka:5.4.3")

  test("do something")(withContainers { container =>
    // 类型转换获取Kafka容器实例
    val kafkaContainer = container.asInstanceOf[KafkaContainer]
    val servers = kafkaContainer.bootstrapServers

    // 创建topic:用容器内执行命令方式
    kafkaContainer.execInContainer(
      "/usr/bin/kafka-topics",
      "--create",
      "--bootstrap-server", "localhost:9092",
      "--replication-factor", "1",
      "--partitions", "1",
      "--topic", "topic1"
    )

    val sparkSession: SparkSession = SparkSession
      .builder()
      .master("local[*]")
      .appName("Unit testing")
      .getOrCreate()

    val df = sparkSession.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", servers)
      .option("subscribe", "topic1")
      .load()

    df.printSchema()
  })
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 04:09:02