"""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.json_schema import build_json_schema_serializer TOPIC = 'event.owsContentReview2.review' SCHEMA_VERSION = 2 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.""" json_schema_value_serializer = build_json_schema_serializer( SCHEMA_REGISTRY_URL, TOPIC, SCHEMA_VERSION ) producer = EventProducer( BOOTSTRAP_SERVERS, StringSerializer(), json_schema_value_serializer, PROTOCOL ) producer.produce( TOPIC, 'example-product-1', mock_product_review_event_json(), delivery_callback ) if __name__ == '__main__': main()