import datetime from functools import cache 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 ( KAFKA_BROKERS, KAFKA_RESYNC_PAYMENT_READINESS_TOPIC, KAFKA_SECURITY_PROTOCOL, ) @cache def get_producer() -> EventProducer: """Get kafka producer.""" return EventProducer( bootstrap_servers=KAFKA_BROKERS, key_serializer=StringSerializer(), value_serializer=SimpleJSONSerializer(), security_protocol=KAFKA_SECURITY_PROTOCOL, ) def emit_resync_payment_readiness_event(account_payee_id: int) -> None: event_key = f'account_payee_id_{account_payee_id}' event_value = { 'timestamp': datetime.datetime.now().isoformat(), 'accountPayee': {'accountPayeeId': f'{account_payee_id}'}, } get_producer().produce( KAFKA_RESYNC_PAYMENT_READINESS_TOPIC, event_key=event_key, event_value=event_value, )