Il commit automatico asincrono degli offset non riesce
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
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.