Source code for core_cdc.processors.base

"""
Base processor interface for CDC (Change Data Capture)
implementations. This module defines the IProcessor abstract
base class that all CDC processors must implement, along
with common CRUD event types.
"""

import json
from abc import abstractmethod
from collections.abc import Iterable, Iterator
from logging import Logger
from typing import Any

from core_mixins.interfaces.factory import IFactory

from core_cdc.base import DdlAction, DdlEvent, EventType, Record
from core_cdc.targets.base import ITarget

CRUD_EVENTS = (
    EventType.INSERT,
    EventType.UPDATE,
    EventType.DELETE,
)


[docs] class IProcessor(IFactory): """Interface for all "Change Data Capture" implementations"""
[docs] def __init__( # pylint: disable=too-many-arguments,too-many-positional-arguments self, service: str, targets: list[ITarget], logger: Logger, events_to_stream: Iterable[EventType] = CRUD_EVENTS, add_event_timestamp: bool = False, ) -> None: """ :param service: The name of the service that is doing the CDC process. :param events_to_stream: By default, insert, update and, delete operations will be streamed. :param logger: Logger used. :param add_event_timestamp: If True, the column event_timestamp will be added when a table is created (only for the targets that execute SQL, see `ITarget.get_add_column_ddl`). The column could be useful for UPSERT or MERGE operations. """ self.service = service self.targets = targets self.events_to_stream = events_to_stream self.add_event_timestamp = add_event_timestamp self.logger = logger
[docs] @classmethod def registration_key(cls) -> str: return cls.__name__
[docs] def execute(self) -> None: """ Execute the CDC processor to read and process events from the stream. This method processes DDL and DML events, forwarding them to configured targets. """ self.logger.info("Reading events from the stream...") for event in self.get_events(): event_type = self.get_event_type(event) if event_type == EventType.DDL_STATEMENT: self._process_ddl_event(event) elif event_type in self.events_to_stream: for target in self.targets: key = target.registration_key() try: records = self.process_dml_event(event) if records: target.save(records) self.logger.info( json.dumps( { "target": key, "number_of_records": len(records), "event_type": event_type, "schema": records[0].schema_name, "table": records[0].table_name, } ) ) except Exception: # pylint: disable=broad-exception-caught self.logger.exception("[%s] -- Error saving records", key) else: try: self.process_event(event) except Exception: # pylint: disable=broad-exception-caught self.logger.exception("Error processing event")
[docs] def _process_ddl_event(self, event: Any) -> None: """ It converts the raw DDL event into a `DdlEvent` and forwards it to the targets with `execute_ddl` enabled. """ targets = [target for target in self.targets if target.execute_ddl] if not targets: return try: ddl_event = self.get_ddl_event(event) except Exception: # pylint: disable=broad-exception-caught self.logger.exception("Error converting the DDL event") return for target in targets: key = target.registration_key() try: query = target.get_ddl_query(ddl_event) if query is None: self.logger.info( "The DDL event %s is not supported by: %s.", ddl_event.action, key, ) continue target.execute(query) self.logger.info("The below query was executed in: %s.", key) self.logger.info(query) if ( self.add_event_timestamp and ddl_event.action == DdlAction.CREATE_TABLE ): self._add_event_timestamp_column(target, ddl_event) except Exception: # pylint: disable=broad-exception-caught self.logger.exception("[%s] -- Error executing DDL", key)
[docs] def _add_event_timestamp_column(self, target: ITarget, ddl_event: DdlEvent) -> None: """It adds the column `event_timestamp` to the table that was just created""" if ddl_event.schema_name is None or ddl_event.table_name is None: self.logger.warning( "The column event_timestamp was not added: the schema or table " "name could not be determined from the DDL event: %s", ddl_event.statement, ) return target.execute( target.get_add_column_ddl( schema=ddl_event.schema_name, table=ddl_event.table_name, column="event_timestamp", type_="bigint", ) )
[docs] @abstractmethod def get_events(self) -> Iterator[Any]: """It returns an iterator with the events to process"""
[docs] @abstractmethod def get_event_type(self, event: Any) -> EventType: """It returns the event type"""
[docs] @abstractmethod def process_dml_event(self, event: Any) -> list[Record]: """It processes the event and return the records to stream"""
[docs] def get_ddl_event(self, event: Any) -> DdlEvent: """ It converts an event classified as `EventType.DDL_STATEMENT` into a `DdlEvent`, the engine independent representation the targets work with. It must be implemented by the processors that produce DDL events. """ raise NotImplementedError( f"{self.registration_key()} does not implement `get_ddl_event`, " "it is required to process DDL events." )
[docs] def process_event(self, event: Any) -> None: """ It should be implemented if another event (apart from DDL and DML) must be processed... """