Il commit automatico asincrono degli offset non riesce

Sep 24 2020

Ho una domanda sul meccanismo di commit automatico di Kafka. Sto usando Spring-Kafka con l'auto-commit abilitato. Come esperimento, ho disconnesso la connessione del mio consumatore a Kafka per 30 secondi mentre il sistema era inattivo (nessun nuovo messaggio nell'argomento, nessun messaggio in elaborazione). Dopo essermi ricollegato ho ricevuto alcuni messaggi in questo modo:

Asynchronous auto-commit of offsets {cs-1915-2553221872080030-0=OffsetAndMetadata{offset=19, leaderEpoch=0, metadata=''}} failed: Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member. This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records.

Primo, non capisco cosa c'è da impegnare? Il sistema era inattivo (tutti i messaggi precedenti erano già stati salvati). In secondo luogo, il tempo di disconnessione era di 30 secondi, molto meno dei 5 minuti (300000 ms) max.poll.interval.ms Terzo, in un errore incontrollato di Kafka ho ricevuto almeno 30.000 messaggi di questo tipo, che sono stati risolti riavviando il processi. Perché sta succedendo?

Sto elencando qui la mia configurazione consumer:

allow.auto.create.topics = true
        auto.commit.interval.ms = 100
        auto.offset.reset = latest
        bootstrap.servers = [kafka1-eu.dev.com:9094, kafka2-eu.dev.com:9094, kafka3-eu.dev.com:9094]
        check.crcs = true
        client.dns.lookup = default
        client.id =
        client.rack =
        connections.max.idle.ms = 540000
        default.api.timeout.ms = 60000
        enable.auto.commit = true
        exclude.internal.topics = true
        fetch.max.bytes = 52428800
        fetch.max.wait.ms = 500
        fetch.min.bytes = 1
        group.id = feature-cs-1915-2553221872080030
        group.instance.id = null
        heartbeat.interval.ms = 3000
        interceptor.classes = []
        internal.leave.group.on.close = true
        isolation.level = read_uncommitted
        key.deserializer = class org.apache.kafka.common.serialization.StringDeserializer
        max.partition.fetch.bytes = 1048576
        max.poll.interval.ms = 300000
        max.poll.records = 500
        metadata.max.age.ms = 300000
        metric.reporters = []
        metrics.num.samples = 2
        metrics.recording.level = INFO
        metrics.sample.window.ms = 30000
        partition.assignment.strategy = [class org.apache.kafka.clients.consumer.RangeAssignor]
        receive.buffer.bytes = 65536
        reconnect.backoff.max.ms = 1000
        reconnect.backoff.ms = 50
        request.timeout.ms = 30000
        retry.backoff.ms = 100
        sasl.client.callback.handler.class = null
        sasl.jaas.config = null
        sasl.kerberos.kinit.cmd = /usr/bin/kinit
        sasl.kerberos.min.time.before.relogin = 60000
        sasl.kerberos.service.name = null
        sasl.kerberos.ticket.renew.jitter = 0.05
        sasl.kerberos.ticket.renew.window.factor = 0.8
        sasl.login.callback.handler.class = null
        sasl.login.class = null
        sasl.login.refresh.buffer.seconds = 300
        sasl.login.refresh.min.period.seconds = 60
        sasl.login.refresh.window.factor = 0.8
        sasl.login.refresh.window.jitter = 0.05
        sasl.mechanism = GSSAPI
        security.protocol = SSL
        send.buffer.bytes = 131072
        session.timeout.ms = 15000
        ssl.cipher.suites = null
        ssl.enabled.protocols = [TLSv1.2, TLSv1.1, TLSv1]
        ssl.endpoint.identification.algorithm = https
        ssl.key.password = [hidden]
        ssl.keymanager.algorithm = SunX509
        ssl.keystore.location = /home/me/feature-2553221872080030.keystore
        ssl.keystore.password = [hidden]
        ssl.keystore.type = JKS
        ssl.protocol = TLS
        ssl.provider = null
        ssl.secure.random.implementation = null
        ssl.trustmanager.algorithm = PKIX
        ssl.truststore.location = /home/me/feature-2553221872080030.truststore
        ssl.truststore.password = [hidden]
        ssl.truststore.type = JKS
        value.deserializer = class org.springframework.kafka.support.serializer.ErrorHandlingDeserializer2

Risposte

1 mike Sep 24 2020 at 14:21

Primo, non capisco cosa c'è da impegnare?

Hai ragione, non c'è nulla di nuovo da impegnare se non fluiscono nuovi dati. Tuttavia, avendo auto.commit abilitato e il tuo consumatore è ancora in esecuzione (anche senza essere in grado di connettersi al broker), il metodo di polling è ancora responsabile dei seguenti passaggi:

  • Recupera i messaggi dalle partizioni assegnate
  • Attivare l'assegnazione della partizione (se necessario)
  • Commit offset se è abilitato il commit con offset automatico

Insieme al tuo intervallo di 100 ms (vedi auto.commit.intervals), il consumatore cerca ancora di eseguire il commit in modo asincrono della posizione di offset (non modificabile) del consumatore.

In secondo luogo, il tempo di disconnessione è stato di 30 secondi, molto inferiore ai 5 minuti (300000 ms) max.poll.interval.ms

Non è il max.poll.interval a causare il ribilanciamento, ma piuttosto la combinazione della tua heartbeat.interval.msimpostazione e del file session.timeout.ms. Il consumatore invia in un thread in background heartbeat basati sull'impostazione dell'intervallo, nel tuo caso 3 secondi. Se il broker non riceve heartbeat prima della scadenza di questo timeout della sessione (nel tuo caso 15 secondi), il broker rimuoverà questo client dal gruppo e avvierà un ribilanciamento.

Una descrizione più dettagliata della configurazione che ho citato è fornita nella documentazione di Kafka su Consumer Configs

Terzo, in un errore incontrollato di Kafka ho ricevuto almeno 30.000 messaggi di questo tipo, che sono stati risolti riavviando il processo. Perché sta succedendo?

Questa sembra essere una combinazione delle prime due domande, in cui non è possibile inviare i battiti cardiaci e il consumatore sta ancora cercando di impegnarsi attraverso il metodo continuo chiamato sondaggio.

Come @GaryRussell ha menzionato nel suo commento, starei attento a usare auto.commit.enablede piuttosto prendere il controllo sulla gestione dell'offset per te stesso.