"""Helpers for publishing invalid messages to the DLQ.""" from typing import Sequence from confluent_kafka.serialization import StringSerializer from kafka_utils.producer.event import EventProducer from lambdacommon.common_config import logger import config from src.models.preference_change import InvalidPreferenceMessage def send_invalid_messages(messages: Sequence[InvalidPreferenceMessage]) -> None: """Send invalid preference change messages to the configured DLQ topic.""" if not messages: return logger.info(f'Sending {len(messages)} invalid message(s) to DLQ topic {config.settings.kafka_dlq_topic}') event_producer_params = dict( bootstrap_servers=config.settings.kafka_bootstrap_servers, key_serializer=StringSerializer(), value_serializer=StringSerializer(), security_protocol=config.settings.kafka_security_protocol ) with EventProducer(**event_producer_params) as producer: for message in messages: headers = { 'errorType': message.error_type.encode(), 'errorMessage': message.error_message.encode(), } producer.produce( topic=config.settings.kafka_dlq_topic, event_key=message.message_key, event_value=message.message_value, headers=headers, auto_flush=False, )