Micronaut集成Kafka首次请求无法获取元数据问题求助
问题概述
使用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

