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

如何让Docker容器中的Scala Kafka生产者连接同VM的Kafka容器

解决Docker容器中Scala Kafka生产者连接同VM Kafka容器的问题

核心问题分析

你的Scala应用在VM本机运行时能连Kafka,是因为localhost:9092指向VM的Kafka端口映射;但在Docker容器里,localhost指的是容器自身的网络栈,而非VM主机。即使使用--network host,如果Kafka的监听器配置没适配,也可能出现连接问题。

方案一:让Scala容器加入Kafka的Docker Compose网络(推荐)

Docker Compose会自动创建一个专属网络,同一网络内的容器可以通过服务名直接通信:

  1. 查看Kafka的Compose网络
    执行命令列出所有Docker网络:

    docker network ls
    

    找到对应Compose项目的网络(通常命名为<你的项目目录名>_default)。

  2. 运行Scala容器时指定网络

    docker run --network <你的Compose网络名称> <你的Scala镜像名>
    
  3. 修改Scala代码的Kafka地址
    将bootstrap.servers改为Kafka容器的内部服务名+端口(对应Compose里PLAINTEXT监听器的配置):

    @main def kafkaTest: Unit =
      val kafkaProps = new Properties()
    
      // 替换为Kafka服务名和内部端口
      kafkaProps.put("bootstrap.servers", "kafka:29092")
      kafkaProps.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
      kafkaProps.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
      kafkaProps.put("linger.ms", "0")
    
      val producer = new KafkaProducer[String, String](kafkaProps)
    
      // 添加回调查看发送结果,方便调试
      producer.send(new ProducerRecord[String, String]("test-topic", "Test message"), (metadata, exception) => {
        if (exception != null) {
          println(s"发送失败: ${exception.getMessage}")
        } else {
          println(s"发送成功,分区: ${metadata.partition()},偏移量: ${metadata.offset()}")
        }
      })
      producer.flush() // 确保消息发送完成再关闭
      producer.close()
    

方案二:修改Kafka监听器配置(适用于独立容器场景)

如果不想让Scala容器加入Compose网络,可以修改Kafka的advertised.listeners,让它对外暴露VM主机的IP地址:

  1. 更新Docker Compose的Kafka配置
    添加新的监听器和端口映射:

    version: '3.8'
    services:
      zookeeper:
        # 保持原有配置不变
      kafka:
        image: confluentinc/cp-kafka:latest
        hostname: broker
        container_name: broker
        depends_on:
          - zookeeper
        ports:
          - "29092:29092"
          - "9092:9092"
          - "9093:9093" # 新增端口映射
          - "9101:9101"
        environment:
          KAFKA_BROKER_ID: 1
          KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181'
          # 新增PLAINTEXT_VM监听器
          KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT,PLAINTEXT_VM:PLAINTEXT
          # 添加VM IP对应的监听器地址(替换<VM的IP地址>为实际值)
          KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092,PLAINTEXT_VM://<VM的IP地址>:9093
          # 保持其他环境变量不变
          KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
          KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
          KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
          KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
          KAFKA_JMX_PORT: 9101
          KAFKA_JMX_HOSTNAME: localhost
    
  2. 重启Kafka容器

    docker-compose down && docker-compose up -d
    
  3. 修改Scala代码的Kafka地址
    将bootstrap.servers改为VM的IP+新增端口:

    kafkaProps.put("bootstrap.servers", "<VM的IP地址>:9093")
    

方案三:排查host网络模式的问题

如果坚持使用--network host仍无法连接,检查以下几点:

  • 确认VM主机的防火墙已开放9092端口
  • 查看Kafka日志,确认PLAINTEXT_HOST监听器已绑定到0.0.0.0而非仅127.0.0.1
  • 在Scala容器内执行ping localhost和telnet localhost 9092,验证网络连通性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 20:42:21