添加upgrade.from配置后Kafka Streams测试遇端口占用错误求助
upgrade.from配置引发集成测试端口占用问题 问题背景
我用Scala开发了一个基于Kafka Streams的应用,集成测试依赖Embedded Kafka Schema Registry。将Kafka Streams从2.5.1升级到3.3.1后,本地运行、单元测试和集成测试都正常,但按照官方升级指南添加upgrade.from配置(值为"2.5.1")后,集成测试开始抛出java.net.BindException: Address already in use错误。
依赖版本对比
升级前依赖
"org.apache.kafka" %% "kafka-streams-scala" % "2.5.1" "org.apache.kafka" % "kafka-clients" % "5.5.1-ccs" "io.confluent" % "kafka-avro-serializer" % "5.5.1" "io.confluent" % "kafka-schema-registry-client" % "5.5.1" "org.apache.kafka" %% "kafka" % "2.5.1" "io.github.embeddedkafka" %% "embedded-kafka-schema-registry" % "5.5.1"
升级后依赖
"org.apache.kafka" %% "kafka-streams-scala" % "3.3.1" "org.apache.kafka" % "kafka-clients" % "7.3.0-ccs" "io.confluent" % "kafka-avro-serializer" % "7.3.0" "io.confluent" % "kafka-schema-registry-client" % "7.3.0" "org.apache.kafka" %% "kafka" % "3.3.1" "io.github.embeddedkafka" %% "embedded-kafka-schema-registry" % "7.3.0"
集成测试初始化代码
class MySpec extends AnyWordSpec with EmbeddedKafkaConfig with EmbeddedKafka { override protected def beforeAll(): Unit = { super.beforeAll() EmbeddedKafka.start() ... } override protected def afterAll(): Unit = { ... EmbeddedKafka.stop() super.afterAll() } }
求助内容
- 该错误产生的原因及解决办法?
- 确认
upgrade.from配置的使用是否正确?
问题分析与解决方案
1. upgrade.from配置使用正确性确认
你的upgrade.from配置是正确的。根据Kafka Streams官方升级规则,从2.5.1版本升级到3.3.1时,设置upgrade.from="2.5.1"符合要求,这个配置会让Streams应用兼容旧版本的状态存储格式,完成平滑升级。
2. 端口占用错误的原因及解决办法
出现Address already in use的核心原因是:添加upgrade.from后,Kafka Streams启动时会触发版本兼容性检查流程,额外初始化状态存储相关组件,导致Embedded Kafka的端口(Broker端口、Schema Registry端口)被重复绑定,或者测试框架的资源清理不及时。
具体解决方向如下:
手动指定Embedded Kafka端口:
重写测试类的embeddedKafkaConfig,固定端口或开启自动端口检测,避免自动分配的端口冲突:override val embeddedKafkaConfig: EmbeddedKafkaConfig = EmbeddedKafkaConfig( kafkaPort = 9093, schemaRegistryPort = 8082, customProperties = Map( "listeners" -> "PLAINTEXT://localhost:9093", "advertised.listeners" -> "PLAINTEXT://localhost:9093" ) )调整资源清理顺序:
在afterAll中先停止Kafka Streams应用实例,再关闭Embedded Kafka,确保所有占用端口的进程都被正确终止:override protected def afterAll(): Unit = { // 先关闭Kafka Streams应用 streams.close(Duration.ofSeconds(10)) // 再停止Embedded Kafka服务 EmbeddedKafka.stop() super.afterAll() }测试中使用临时状态存储:
给测试用的Streams配置指定临时状态目录,避免兼容性检查时重复加载本地持久化的状态文件,减少端口占用风险:val streamsConfig = new Properties() streamsConfig.put(StreamsConfig.STATE_DIR_CONFIG, Files.createTempDirectory("streams-test").toAbsolutePath.toString)拆分Embedded服务启动逻辑:
分开启动Kafka Broker和Schema Registry,确保服务启动顺序正确,避免端口绑定冲突:override protected def beforeAll(): Unit = { EmbeddedKafka.startKafka() EmbeddedKafka.startSchemaRegistry() super.beforeAll() }
内容的提问来源于stack exchange,提问作者katmhicks

