"""Example producer.""" from confluent_kafka.serialization import StringSerializer from helpers import delivery_callback from helpers import mock_product_review_event_json from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.avro import build_avro_serializer TOPIC = 'event.owsContentReview.review' SCHEMA_VERSION = 5 SCHEMA_REGISTRY_URL = 'https://dev-schema-registry.dev.theorchard.io' BOOTSTRAP_SERVERS = ( '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,' 'b-1.dev-managed-kafka-cdc.rk4es0.c11.kafka.us-east-1.amazonaws.com:9094' ) PROTOCOL = 'SSL' # or PLAINTEXT if local non-ssl def main(): """Produce Message.""" avro_value_serializer = build_avro_serializer(SCHEMA_REGISTRY_URL, TOPIC, SCHEMA_VERSION) producer = EventProducer(BOOTSTRAP_SERVERS, StringSerializer(), avro_value_serializer, PROTOCOL) producer.produce( TOPIC, 'example-product-1', mock_product_review_event_json(), delivery_callback ) if __name__ == '__main__': main()