aio-libs / aio-libs/aiokafka

[QUESTION]

Aberta
#821 0 comentários 0 reações 0 responsáveis Ver no GitHub
question
Linguagem predominante
Python
Estrelas
1.4k
Forks
269
Merge médio
1d 1h
PRs com merge (30d)
6

Descrição

I am using Python 3.9.9. I am facing the problem of message duplicate even after a successful commit. The same message appears multiple times even after committing. I can see very frequent logs. I have filtered my log and seen 10 produced message appears around 16 times committed of the same offset. I am not sure where to start looking at.

I need guidance to fix the issue.

**
WARNING -- Heartbeat failed: local member_id was not recognized; resetting and re-joining group

ERROR -- Heartbeat session expired - marking coordinator dead

WARNING -- Heartbeat failed for group EL_API_Cluster because it is rebalancing
INFO -- Revoking previously assigned partitions frozenset({TopicPartition(topic='neo-transfo
rm-topic', partition=0)}) for group EL_API_Cluster
**

Parameter settings
```
AEH_CONSUMER_CONFIG = dict(
bootstrap_servers=os.getenv('ENVIRONMENT_VARIABLE_AEH_SERVER_PORT').strip(),
group_id=os.getenv('ENVIRONMENT_VARIABLE_AEH_GROUP_ID').strip(),
sasl_plain_username=os.getenv('ENVIRONMENT_VARIABLE_AEH_USERNAME').strip(),
sasl_plain_password=get_secrets().sasl_pw.strip(),
enable_auto_commit=False,
sasl_mechanism='PLAIN',
security_protocol='SASL_SSL',
client_id='python-consumer-client',
rebalance_timeout_ms=300000,
max_poll_records=10,
ca_path=None)
```

Consumer Code

```
import asyncio
import copy
import datetime
import json
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer
from aiokafka.helpers import create_ssl_context
from aiokafka.errors import KafkaError, ConsumerStoppedError, RecordTooLargeError, InvalidMessageError
from kafka import TopicPartition, OffsetAndMetadata
from kafka.errors import KafkaTimeoutError

from source.settings import logger, STATUS_MESSAGE, LIFE_CYCLE_PUBLISH_TOPIC
from source.util import Singleton

class ContextSSL:
"""base class to consumer producer"""

def __init__(self, ca_path=None):
self.ssl_context = create_ssl_context()
if ca_path:
self.ssl_context = create_ssl_context(cadata=ca_path)

class ConsumerHandler:
"""Consume message asynchronously"""

def __init__(self, topic, **conf):
context = ContextSSL(conf.pop('ca_path'))
self.consumer = AIOKafkaConsumer(**conf, ssl_context=context.ssl_context)
self._topic = topic

async def get_message(self, callback):
"""listen to message"""
try:
await self.consumer.start()
self.consumer.subscribe(pattern=self._topic)
logger.info(f'Clustor topic {await self.consumer.topics()}')
async for msg in self.consumer:
try:
logger.info(
f"consumed: topic: {msg.topic}, partition: {msg.partition}, offset: {msg.offset}, "
f"key: {msg.key}, body: {msg.value}, timestamp:{msg.timestamp} header: {msg.headers}")
asyncio.get_running_loop().create_task(callback(message=msg.value, headers=msg.headers))
await self.commit_on_offset(partition=msg.partition, topic=msg.topic, offset=msg.offset + 1)
except RecordTooLargeError as record_error:
logger.critical(f"Failed to consume message from consumer because max_partition_fetch_bytes is "
f"low: {record_error}")
except InvalidMessageError as message_error:
logger.critical(f'CRC check on MessageSet failed due to connection failure or bug. '
f'Always raised. Changed in version 0.5.0, before we ignored this error in '
f'async for. {message_error}')
except KafkaError as err:
logger.critical(f"Failed to consume message from consumer broker: {err}")
except KafkaError as err:
logger.critical(f"Failed to start consumer: {err}")
except ConsumerStoppedError as err:
logger.critical(f"Unexpectedly consumer closed: {err}")
finally:
await self.consumer.stop()

async def commit_on_offset(self, partition, topic, offset):
"""commit to offset"""
topic_partition = TopicPartition(topic=topic, partition=partition)
offset_metadata = OffsetAndMetadata(offset=offset, metadata='')
offset_meta_dict = dict()
offset_meta_dict[topic_partition] = offset_metadata
await self.consumer.commit(offsets=offset_meta_dict)
last_committed_offset = await self.consumer.committed(partition=topic_partition)
logger.info(f'offset: {last_committed_offset} committed on partition: {partition} at topic: {topic}')
```

Thanks

Guia de contribuição

Abrir o guia de contribuição

Avaliação

Esta issue ainda não foi avaliada.

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.