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

添加upgrade.from配置后Kafka Streams测试遇端口占用错误求助

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()
  }
}

求助内容

  1. 该错误产生的原因及解决办法?
  2. 确认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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 12:50:23