You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Describe the bug
I was able to connect to Azure Event Hub but when I run producer and consumer and connect to a topic I see an Incoming (Sum) and Successful Request (Sum) spike to 4K and stays constantly at 4K even when no message is sent
Expected behavior
As there is only 1 pod running producer and consumer there should be fewer Incoming and Successful Request (Sum) because I have one more Kubernetes which is running producer and consumer written in Node.js which hardly consumes 1k Incoming and Successful Request (Sum) which has a minimum of 5 Pods running all the time
Environment (please complete the following information):
aiokafka version (python -c "import aiokafka; print(aiokafka.__version__)"): "0.8.1"
Kafka Broker version (kafka-topics.sh --version):
Other information (Confluent Cloud version, etc.): Azure Event Hub as Kafka
Reproducible example
# Add a short Python script or Docker configuration that can reproduce the issue.Consumerasyncdef_subscribe_to_topic(self, topic: str, group_id: str):
consumer=AIOKafkaConsumer(topic,
bootstrap_servers=self.server,
group_id=group_id,
sasl_plain_username='$ConnectionString',
sasl_plain_password=self.kafkaPass,
sasl_mechanism='PLAIN',)
awaitconsumer.start()
print("consumer started")
try:
asyncformsginconsumer:
payload=json.loads(msg.value)
print(payload)
try:
charging_station_id=payload.get("charging_station_id")
awaitself._process_payload(charging_station_id, payload)
exceptExceptionase:
logger.error(f"Error processing payload: {e}, payload: {payload}")
finally:
awaitconsumer.stop()
print("consumer stopped")
Producerasyncdef_publish(self, topic: str, envelope: Envelope):
envelope.timestamp=datetime.now(timezone.utc)
endpoint=os.getenv("PUPSUB_KAFKA_SERVER")
kafkaPass=os.getenv("PUPSUB_KAFKA_PASSWORD")
ifendpointisNoneorkafkaPassisNone:
raiseValueError("PUPSUB_KAFKA_ENDPOINT not set or PUPSUB_KAFKA_PASSWORD not set ")
producer=AIOKafkaProducer(bootstrap_servers=endpoint,
sasl_mechanism='PLAIN',
sasl_plain_username='$ConnectionString',
sasl_plain_password=kafkaPass,
enable_idempotence=True)
# Get cluster layout and initial topic/partition leadership informationawaitproducer.start()
try:
# Produce messageawaitproducer.send_and_wait(topic, envelope.json())
finally:
# Wait for all pending messages to be delivered or expire.awaitproducer.stop()
The text was updated successfully, but these errors were encountered:
Describe the bug
I was able to connect to Azure Event Hub but when I run producer and consumer and connect to a topic I see an Incoming (Sum) and Successful Request (Sum) spike to 4K and stays constantly at 4K even when no message is sent
Expected behavior
As there is only 1 pod running producer and consumer there should be fewer Incoming and Successful Request (Sum) because I have one more Kubernetes which is running producer and consumer written in Node.js which hardly consumes 1k Incoming and Successful Request (Sum) which has a minimum of 5 Pods running all the time
Environment (please complete the following information):
python -c "import aiokafka; print(aiokafka.__version__)"
): "0.8.1"kafka-topics.sh --version
):Reproducible example
The text was updated successfully, but these errors were encountered: