yiorgos icon

confluent-kafka-producer

yiorgos | PRO | 03/26/21 04:08:47 PM UTC (Edited) | 0 ⭐ | 280 👁️ | Never ⏰ | []
Python |

1.03 KB

|

None

|

0 👍

/

0 👎

# Simple Kafka Producer, using confluent-kafka-python (librdkafka)
 
# Sample dockerfile to run this producer
#
# FROM python:3
# WORKDIR /usr/src/app
# RUN pip install --no-cache-dir  confluent-kafka==1.0.0
# COPY ca-cert .
# COPY rd.py .
# CMD ["python", "./rd.py"]
 
from confluent_kafka import Producer
 
def delivery_report(err, msg):
  if err is not None:
    print("Message delivery failed: {}".format(err))
  else:
    print('Message delivered to {} [{}]'.format(msg.topic(), msg.partition()))
 
topic = 'CLIENT_NAME-tests'
 
p = Producer({
  'bootstrap.servers': 'server.name:9094',
  'security.protocol': 'SASL_SSL',
  'ssl.ca.location': 'ca-cert',
  'sasl.mechanisms': 'PLAIN',
  'sasl.username': 'USERNAME',
  'sasl.password': 'XXX_YYY_ZZZ',
  'acks': 1,
  'compression.type': 'lz4',
  })
 
for i in range(1000000):
  message = "Event {}".format(i)
  p.produce(topic, message, on_delivery=delivery_report)
  print("{} Queue len: {}".format(message, len(p)))
  if len(p) > 90000:
    p.poll(0)
 
p.flush()
 

Comments