from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from kafka_utils.producer.serializer.simple_json import SimpleJSONSerializer from config import config producer = EventProducer( config.KAFKA_BROKERS, StringSerializer(), SimpleJSONSerializer(), 'SSL') def mycallback(e, f): pass def send_kafka_message(contribution_id): payload = { 'contributionId': contribution_id, 'application': 'neighbouring-rights-migration', 'source': 'backfill' } producer.produce( config.CONTRIB_METADATA_TOPIC, contribution_id, payload, mycallback )