"""
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 _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,
}