"""
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...
"""