Source code for core_cdc.base

"""
Base classes and types for CDC (Change Data Capture)
operations. This module defines the core EventType and
DdlAction enums and the Record and DdlEvent classes used
across all CDC processors to represent database change
events (data and schema changes).
"""

from __future__ import annotations

import json
import uuid
from dataclasses import dataclass, field
from datetime import date, time, timedelta
from decimal import Decimal
from typing import Any

from core_mixins.compatibility import StrEnum


[docs] def _format_time_delta(value: timedelta) -> str: """ It formats a MySQL TIME (that has hours over 24 and negative values, so it can not be a `datetime.time`) as MySQL does: [-]HH:MM:SS[.ffffff] """ microseconds = (value.days * 86400 + value.seconds) * 1_000_000 + value.microseconds sign = "-" if microseconds < 0 else "" seconds, microseconds = divmod(abs(microseconds), 1_000_000) hours, seconds = divmod(seconds, 3600) minutes, seconds = divmod(seconds, 60) text = f"{sign}{hours:02d}:{minutes:02d}:{seconds:02d}" return f"{text}.{microseconds:06d}" if microseconds else text
[docs] def _to_json_value(value: Any) -> Any: # pylint: disable=too-many-return-statements """ It converts the values of a record that JSON does not support: dates and times are ISO formatted, decimals and UUIDs are strings (a decimal as a float would lose precision), bytes are hexadecimal and sets are sorted lists. The rest of the values (including the unknown types) are not changed. """ if isinstance(value, (date, time)): # datetime is a date... return value.isoformat() if isinstance(value, timedelta): return _format_time_delta(value) if isinstance(value, (Decimal, uuid.UUID)): return str(value) if isinstance(value, (bytes, bytearray)): return value.hex() if isinstance(value, (set, frozenset)): return sorted((_to_json_value(item) for item in value), key=str) if isinstance(value, dict): return {key: _to_json_value(item) for key, item in value.items()} if isinstance(value, (list, tuple)): return [_to_json_value(item) for item in value] return value
[docs] class EventType(StrEnum): """Enumeration of CDC event types.""" GLOBAL = "GLOBAL" # Event related to the processor itself. DDL_STATEMENT = "QUERY" # Event for DDL statement (like create table, etc.). INSERT = "INSERT" # Event for DML INSERT operation. UPDATE = "UPDATE" # Event for DML UPDATE operation. DELETE = "DELETE" # Event for DML DELETE operation.
[docs] class DdlAction(StrEnum): """ Kind of schema change described by a `DdlEvent`, no matter the engine it comes from. "Schema" is a database and "table" a collection when the source is MongoDB. """ CREATE_SCHEMA = "CREATE_SCHEMA" DROP_SCHEMA = "DROP_SCHEMA" CREATE_TABLE = "CREATE_TABLE" ALTER_TABLE = ( "ALTER_TABLE" # Also changes of the collection (rename, indexes, etc.). ) DROP_TABLE = "DROP_TABLE" OTHER = ( "OTHER" # Anything else (create index, truncate, sharding operations, etc.). )
[docs] @dataclass class Record: # pylint: disable=too-many-instance-attributes """ It provides a wrapper or common object useful for integration across services that needs to handle data replication via Change Data Capture producers and consumers... :param event_type: It specifies the operation that generated the event. :param record: The record (like a record/row in a database). :param service: The service from where the record came. :param schema_name: Schema from where the record was retrieved. :param table_name: Table from where the record was retrieved. :param primary_key: Attributes to use to identify a record as unique. :param global_id: Unique identifier (e.g. Global Transaction Identifier in MySQL). :param transaction_id: Identifier for the transaction (e.g. `<binlog file>:<position of BEGIN>` in MySQL, `<session UUID>:<txnNumber>` in MongoDB). :param event_timestamp: Timestamp when the record was generated. :param source: The source from where the record was retrieved (e.g. binlog file name). :param position: Record position in the source. """ event_type: EventType record: dict service: str schema_name: str table_name: str primary_key: tuple[Any, ...] | str global_id: str | None = None transaction_id: str | None = None event_timestamp: int | None = None source: str | None = None position: int | None = None def __str__(self) -> str: return json.dumps(self.to_json())
[docs] def to_json(self) -> dict: """ It returns the JSON version of the record required to be streamed... The values of the record that JSON does not support are converted: dates and times are ISO formatted, decimals and UUIDs are strings (a float would lose precision), bytes are hexadecimal, a MySQL TIME is `[-]HH:MM:SS[.ffffff]` and a set is a sorted list. It does not change the record and the unknown types are not converted. :return: A dictionary that follows the below structure: Example:: { "global_id": "GTID in MySQL | None", "transaction_id": "binlog.000001:1560 (MySQL) | uuid:txn (MongoDB) | None", "event_timestamp": 1653685384, "event_type": "INSERT | UPDATE | DELETE", "service": "service-name", "source": "binlog.000001 | FileName | Something", "position": 8077, "primary_key": "(str | tuple)", "schema_name": "schema_name", "table_name": "table_name", "record": { ... // This is an example... "id": "000-1", "category": "Marketing", "price": 100.0, ... } } """ return { "global_id": self.global_id, "transaction_id": self.transaction_id, "event_timestamp": self.event_timestamp, "event_type": self.event_type.value, "service": self.service, "source": self.source, "position": self.position, "primary_key": self.primary_key, "schema_name": self.schema_name, "table_name": self.table_name, "record": { key: _to_json_value(value) for key, value in self.record.items() }, }
[docs] @dataclass class DdlEvent: # pylint: disable=too-many-instance-attributes """ It is the DDL (schema change) counterpart of `Record`: the processors convert the raw DDL event of their engine into it, so the targets do not depend on the engine the change comes from. :param action: The kind of schema change. :param service: The service from where the event came. :param schema_name: Schema (database) affected. It could be None if it is unknown. :param table_name: Table (collection) affected. None for schema level changes. :param statement: The statement as the source engine wrote it (SQL in MySQL), or None when the engine does not have one (like MongoDB). :param global_id: Unique identifier (e.g. Global Transaction Identifier in MySQL). :param event_timestamp: Timestamp when the change was generated. :param source: The source from where the event was retrieved (e.g. binlog file name). :param position: Event position in the source. :param raw: The original event of the engine, for the targets that need more details than the ones above. It is not included in `to_json()`. """ action: DdlAction service: str schema_name: str | None = None table_name: str | None = None statement: str | None = None global_id: str | None = None event_timestamp: int | None = None source: str | None = None position: int | None = None raw: Any = field(default=None, repr=False, compare=False) def __str__(self) -> str: return json.dumps(self.to_json())
[docs] def to_json(self) -> dict: """ It returns the JSON version of the event. :return: A dictionary that follows the below structure: Example:: { "action": "CREATE_TABLE", "service": "service-name", "schema_name": "schema_name", "table_name": "table_name", "statement": "CREATE TABLE schema_name.table_name (...)", "global_id": "GTID in MySQL | None", "event_timestamp": 1653685384, "source": "binlog.000001 | None", "position": 8077 } """ return { "action": self.action.value, "service": self.service, "schema_name": self.schema_name, "table_name": self.table_name, "statement": self.statement, "global_id": self.global_id, "event_timestamp": self.event_timestamp, "source": self.source, "position": self.position, }