"""Kafka Outbox Iterable.""" from __future__ import annotations import base64 import json from typing import Iterator from lambdacommon.common_config import logger from src.enums import DebeziumOperation from src.schemas import AbacusOutbox from src.schemas.kafka import DebeziumCDCEvent, KafkaEvent, KafkaRecord class KafkaOutboxIterable: """Kafka/MSK CDC event iterable. Iterates over AbacusOutbox events extracted from Debezium CDC (Change Data Capture) events received from Kafka/MSK. Only processes INSERT operations (CREATE). Event Processing: - Filters for CREATE operations only (inserts to outbox table) - Skips UPDATE, DELETE, and snapshot READ operations - Extracts outbox events from CDC 'after' payload - Continues iteration on individual record failures Idempotency: Works in conjunction with OutboxProcessor's idempotency checks to handle Kafka redelivery scenarios. The processor skips already-completed events. """ def __init__(self, kafka_event: KafkaEvent): """Initialize Kafka outbox iterable. Args: kafka_event: Kafka event containing Debezium CDC records from Lambda's Kafka/MSK event source mapping """ self._kafka_event = kafka_event def __iter__(self) -> Iterator[AbacusOutbox]: """Iterate over outbox events from Kafka CDC records. Processes all partitions in the Kafka event batch. For each CDC record, extracts the outbox event if it's a CREATE operation. Skips invalid or non-CREATE records without failing the entire batch. Yields: AbacusOutbox: Parsed outbox events from CREATE operations Note: Only CREATE (INSERT) operations are yielded. All other Debezium operations (UPDATE, DELETE, READ) are silently skipped. """ for topic_partition, partition_records in self._kafka_event.records.items(): logger.info(f'Processing from partition: {topic_partition}') for record in partition_records: cdc_event = self._extract_cdc_event(record) if cdc_event and cdc_event.op == DebeziumOperation.CREATE: outbox_event = self._extract_outbox_event(cdc_event) if outbox_event: yield outbox_event def _extract_cdc_event(self, record: KafkaRecord) -> DebeziumCDCEvent | None: """Extract and parse Debezium CDC event from Kafka record. Args: record: Kafka record with base64-encoded CDC event. Returns: DebeziumCDCEvent: Parsed CDC event, or None if extraction fails. """ try: value_json = base64.b64decode(record.value).decode('utf-8') return DebeziumCDCEvent(**json.loads(value_json)) except (ValueError, json.JSONDecodeError, KeyError) as e: logger.error(f'Failed to extract CDC event: {e}') return None def _extract_outbox_event(self, cdc_event: DebeziumCDCEvent) -> AbacusOutbox | None: """Extract AbacusOutbox event from CDC event 'after' payload. Args: cdc_event: Debezium CDC event containing the change data. Returns: AbacusOutbox: Parsed outbox event, or None if extraction fails. """ try: if cdc_event.after: return AbacusOutbox(**cdc_event.after) logger.warning('CDC event missing "after" payload') return None except (ValueError, json.JSONDecodeError, KeyError) as e: logger.error(f'Failed to extract Outbox event: {e}') return None