Interfaces#
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).
- core_cdc.base._format_time_delta(value: timedelta) str[source]#
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]
- core_cdc.base._to_json_value(value: Any) Any[source]#
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.
- class core_cdc.base.EventType(*values)[source]#
Bases:
StrEnumEnumeration of CDC event types.
- GLOBAL = 'GLOBAL'#
- DDL_STATEMENT = 'QUERY'#
- INSERT = 'INSERT'#
- UPDATE = 'UPDATE'#
- DELETE = 'DELETE'#
- static _generate_next_value_(name, start, count, last_values)#
Return the lower-cased version of the member name.
- class core_cdc.base.DdlAction(*values)[source]#
Bases:
StrEnumKind 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'#
- DROP_TABLE = 'DROP_TABLE'#
- OTHER = 'OTHER'#
- static _generate_next_value_(name, start, count, last_values)#
Return the lower-cased version of the member name.
- class core_cdc.base.Record(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)[source]#
Bases:
objectIt provides a wrapper or common object useful for integration across services that needs to handle data replication via Change Data Capture producers and consumers…
- Parameters:
event_type – It specifies the operation that generated the event.
record – The record (like a record/row in a database).
service – The service from where the record came.
schema_name – Schema from where the record was retrieved.
table_name – Table from where the record was retrieved.
primary_key – Attributes to use to identify a record as unique.
global_id – Unique identifier (e.g. Global Transaction Identifier in MySQL).
transaction_id – Identifier for the transaction (e.g. <binlog file>:<position of BEGIN> in MySQL, <session UUID>:<txnNumber> in MongoDB).
event_timestamp – Timestamp when the record was generated.
source – The source from where the record was retrieved (e.g. binlog file name).
position – Record position in the source.
- to_json() dict[source]#
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.
- Returns:
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, ... } }
- class core_cdc.base.DdlEvent(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 = None)[source]#
Bases:
objectIt 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.
- Parameters:
action – The kind of schema change.
service – The service from where the event came.
schema_name – Schema (database) affected. It could be None if it is unknown.
table_name – Table (collection) affected. None for schema level changes.
statement – The statement as the source engine wrote it (SQL in MySQL), or None when the engine does not have one (like MongoDB).
global_id – Unique identifier (e.g. Global Transaction Identifier in MySQL).
event_timestamp – Timestamp when the change was generated.
source – The source from where the event was retrieved (e.g. binlog file name).
position – Event position in the source.
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().
- to_json() dict[source]#
It returns the JSON version of the event.
- Returns:
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 }