from signal import signal from signal import SIGINT from sys import exit from threading import Event from threading import Thread from uuid import uuid4 import random import time from confluent_kafka import DeserializingConsumer from confluent_kafka import SerializingProducer from confluent_kafka.schema_registry import SchemaRegistryClient from confluent_kafka.schema_registry.avro import AvroDeserializer from confluent_kafka.schema_registry.avro import AvroSerializer from faker import Faker from icecream import ic import certifi MSK_CONFIG = { 'bootstrap.servers': 'b-1.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094,' 'b-2.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094,' 'b-3.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094', 'security.protocol': 'SSL', 'ssl.ca.location': certifi.where() } SCHEMA_REGISTRY_CONFIG = { 'url': 'https://dev-schema-registry.dev.theorchard.io' } MESSAGE_SCHEMA = """ { "namespace": "com.orchard.kafka.data.highway", "name": "TestAvroArtist", "type": "record", "fields": [ {"name": "first_name", "type": "string"}, {"name": "last_name", "type": "string"}, {"name": "age", "type": "int"} ] } """ schema_registry_client = SchemaRegistryClient(SCHEMA_REGISTRY_CONFIG) # deserializer = AvroDeserializer( # schema_registry_client=schema_registry_client) serializer = AvroSerializer( schema_registry_client=schema_registry_client, schema_str=MESSAGE_SCHEMA) TOPIC = 'topic.aminakov.local.avro' # your topic should have >= 2 replicas fake = Faker() do_work_event = Event() def acked(err, msg): if err is not None: ic(err) else: print('Produced record to topic {} partition [{}] @ offset {}'. format(msg.topic(), msg.partition(), msg.offset())) def produce(producer): global do_work_event while do_work_event.is_set(): msg = { 'first_name': fake.first_name(), 'last_name': fake.last_name(), 'age': random.randint(10, 100) } key = msg['first_name'] + msg['last_name'] + str(msg['age']) producer.produce( topic=TOPIC, key=key, value=msg, on_delivery=acked, headers={'CorrelationId': str(uuid4())} ) producer.poll(0) time.sleep(2) ic(producer.flush(2)) ic('Producer stopped.') def consume(consumer): global do_work_event consumer.subscribe([TOPIC]) while do_work_event.is_set(): msg = consumer.poll(10) if msg and not msg.error(): ic(msg.key()) ic(msg.value()) consumer.close() ic('Consumer stopped.') def signal_handler(sig, frame): global do_work_event ic('Exitting due to signal..') do_work_event.clear() exit(0) if __name__ == '__main__': signal(SIGINT, signal_handler) producer_config = MSK_CONFIG.copy() producer_config.update({ 'value.serializer': serializer }) producer = SerializingProducer(producer_config) # consumer_config = MSK_CONFIG.copy() # consumer_config.update({ # 'value.deserializer': deserializer, # 'group.id': 'aminakov_local_avro_consumer', # 'auto.offset.reset': 'earliest' # }) # consumer = DeserializingConsumer(consumer_config) do_work_event.set() producer_thread = Thread(target=produce, args=[producer]) #consumer_thread = Thread(target=consume, args=[consumer]) producer_thread.start() #consumer_thread.start() producer_thread.join() #consumer_thread.join()