使用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
相关产品推荐
相关产品推荐

