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

Micronaut集成Kafka首次请求无法获取元数据问题求助

解决Kafka首次发送消息LEADER_NOT_AVAILABLE警告及消息零丢失保障

问题概述

使用Micronaut 3.6.3应用向Docker Compose部署的Kafka发送消息时,首次请求出现以下警告:

[Producer clientId=producer-1] Error while fetching metadata with correlation id 1 : {accountRegistered=LEADER_NOT_AVAILABLE}

后续消息发送恢复正常,但需要确保账户注册类消息零丢失。

核心原因

该警告本质是Kafka broker初始化未完成,且accountRegistered主题为首次访问时动态创建,此时主题的leader副本尚未完成选举,导致生产者首次请求失败。虽然后续Kafka会自动完成主题创建和leader选举,但如果生产者未配置可靠的重试策略,首次发送的消息可能丢失。

解决方案

1. 提前预创建目标主题

避免动态创建主题的延迟,在Kafka启动阶段就创建accountRegistered主题,确保应用连接时主题已就绪。

修改Docker Compose的Kafka服务配置,增加主题创建命令:

services:
  kafka:
    image: 'bitnami/kafka:3.2'
    hostname: 'kafka'
    environment:
      ALLOW_PLAINTEXT_LISTENER: 'yes'
      KAFKA_BROKER_ID: 1
      KAFKA_CFG_ADVERTISED_LISTENERS: 'INSIDE://kafka:29092, OUTSIDE://localhost:9092'
      KAFKA_CFG_INTER_BROKER_LISTENER_NAME: 'INSIDE'
      KAFKA_CFG_LISTENERS: 'INSIDE://:29092, OUTSIDE://:9092'
      KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: 'INSIDE:PLAINTEXT, OUTSIDE:PLAINTEXT'
      KAFKA_CFG_ZOOKEEPER_CONNECT: 'zookeeper:2181'
    ports:
      - '9092:9092'
    depends_on:
      - 'zookeeper'
    # 新增:启动后创建主题
    command: >
      bash -c "
      /opt/bitnami/scripts/kafka/run.sh &
      sleep 10 &&
      /opt/bitnami/kafka/bin/kafka-topics.sh --create --topic accountRegistered --bootstrap-server kafka:29092 --partitions 3 --replication-factor 1
      wait
      "
  zookeeper:
    image: 'bitnami/zookeeper:3.8'
    hostname: 'zookeeper'
    environment:
      ALLOW_ANONYMOUS_LOGIN: 'yes'
    ports:
      - '2181:2181'

2. 配置生产者可靠性参数

在Micronaut的application.yml中强化生产者配置,确保消息发送失败时自动重试,且保证消息不重复、不丢失:

kafka:
  bootstrap:
    servers: 'localhost:9092'
  producers:
    default:
      retries: 10                  # 最大重试次数
      acks: all                    # 等待所有同步副本确认消息接收
      retry-backoff-ms: 1000       # 重试间隔时间
      enable-idempotence: true     # 开启幂等性,避免重试导致重复消息
      delivery-timeout-ms: 300000  # 最大投递超时(5分钟)

3. 优化Docker容器启动健康检查

仅靠depends_on无法确保Kafka/ZooKeeper真正就绪,添加健康检查保证服务完全可用后再允许应用连接:

services:
  zookeeper:
    image: 'bitnami/zookeeper:3.8'
    hostname: 'zookeeper'
    environment:
      ALLOW_ANONYMOUS_LOGIN: 'yes'
    ports:
      - '2181:2181'
    healthcheck:
      test: ["CMD", "zkServer.sh", "status"]
      interval: 10s
      timeout: 5s
      retries: 5
  kafka:
    image: 'bitnami/kafka:3.2'
    hostname: 'kafka'
    environment:
      ALLOW_PLAINTEXT_LISTENER: 'yes'
      KAFKA_BROKER_ID: 1
      KAFKA_CFG_ADVERTISED_LISTENERS: 'INSIDE://kafka:29092, OUTSIDE://localhost:9092'
      KAFKA_CFG_INTER_BROKER_LISTENER_NAME: 'INSIDE'
      KAFKA_CFG_LISTENERS: 'INSIDE://:29092, OUTSIDE://:9092'
      KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: 'INSIDE:PLAINTEXT, OUTSIDE:PLAINTEXT'
      KAFKA_CFG_ZOOKEEPER_CONNECT: 'zookeeper:2181'
    ports:
      - '9092:9092'
    depends_on:
      zookeeper:
        condition: service_healthy  # 等待ZooKeeper健康检查通过
    healthcheck:
      test: ["CMD", "/opt/bitnami/kafka/bin/kafka-topics.sh", "--list", "--bootstrap-server", "localhost:9092"]
      interval: 10s
      timeout: 5s
      retries: 10

4. 消费者端强化消息处理可靠性

现有消费者已设置offsetReset = OffsetReset.EARLIEST,可确保重启后消费最早未处理的消息。进一步改为手动提交offset,确保消息处理完成后才确认:

@KafkaListener(offsetReset = OffsetReset.EARLIEST, ackMode = AckMode.MANUAL)
class AccountReferenceUpdaterEventConsumer {

    @Inject
    AccountReferenceEntityRepository accountReferenceEntityRepository

    @Topic('accountRegistered')
    void receive(@MessageBody AccountRegisteredEvent event, Acknowledgement acknowledgement) {

        def account = event.source
        accountReferenceEntityRepository.findById(account.id)
                .ifPresentOrElse(
                        accountReference -> log.warn('Account {} already registered', account.id),
                        () -> {
                            def accountReference = new AccountReferenceEntity(
                                    accountId: account.id,
                                    username: account.username
                            )
                            accountReferenceEntityRepository.save(accountReference)
                        }
                )
        // 消息处理完成后手动提交offset
        acknowledgement.acknowledge()
    }
}

效果验证

完成以上配置后,应用启动时Kafka主题已就绪,生产者首次发送不会出现LEADER_NOT_AVAILABLE警告;同时通过生产者重试、副本确认、消费者手动提交等机制,确保账户注册消息全程不丢失。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 12:43:28