Micronaut with kafka fails to retrieve metadata on first attempts

Viewed 33

I'm trying to send messages from a Micronaut 3.6.3 application to Kafka deployed with docker-compose. On first attempt I receive a warning like this:

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

For the following messages, the problem disappear but my requirement is to not lost any about account registration.

My docker compose configuration:

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'
  #TODO: Can be removed with the future versions of Kafka (using KRaft)
  zookeeper:
    image: 'bitnami/zookeeper:3.8'
    hostname: 'zookeeper'
    environment:
      ALLOW_ANONYMOUS_LOGIN: 'yes'
    ports:
      - '2181:2181'

From the application I use 'localhost:9092' to connect.

My consumer code:

@KafkaListener(offsetReset = OffsetReset.EARLIEST)
class AccountReferenceUpdaterEventConsumer {

    @Inject
    AccountReferenceEntityRepository accountReferenceEntityRepository

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

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

}

application.yml:

kafka:
  bootstrap:
    servers: 'localhost:9092'
0 Answers
Related