MySQL#

MySQL BinLog Processor#

Tested with MySQL 9.7, 8.4, 8.0 and 5.7, see Tested Versions in the index.

MySQL BinLog processor for Change Data Capture services.

core_cdc.processors.mysql.mysql_binlog._clean(statement: str) → str[source]#

It removes the comments of a statement and converts it to a single line

core_cdc.processors.mysql.mysql_binlog._decode_json(value: Any) → Any[source]#

mysql-replication returns the keys and the strings of a JSON column as bytes (UTF-8), it converts them to text. Any other value is not changed.

core_cdc.processors.mysql.mysql_binlog._parse_values(column_type: str) → list[str][source]#

The values of an ENUM or SET column type: enum(‘a’,’it’’s’) -> [“a”, “it’s”]

core_cdc.processors.mysql.mysql_binlog._encoding_name(charset: str) → str | None[source]#

The name of a MySQL character set that mysql-replication can decode, if any

core_cdc.processors.mysql.mysql_binlog._unquote(identifier: str) → str[source]#

It removes the backticks of an identifier (a``b -> a`b)

class core_cdc.processors.mysql.mysql_binlog.MySqlBinlogProcessor(stream: BinLogStreamReader, connection_settings: dict | None = None, **kwargs)[source]#

Bases: IProcessor

It processes the events from the BinLog files.

The binary log contains “events” that describe database changes such as table creation operations or changes to table data. It also contains events for statements that potentially could have made changes (for example, a DELETE which matched no rows), unless row-based logging is used. The binary log also contains information about how long each statement took that updated data.

More information: https://dev.mysql.com/doc/refman/8.0/en/binary-log.html

__init__(stream: BinLogStreamReader, connection_settings: dict | None = None, **kwargs) → None[source]#

https://python-mysql-replication.readthedocs.io/en/stable/binlogstream.html :param stream: BinLogStreamReader object.

get_events() → Iterator[Any][source]#

It returns an iterator with the events to process

get_event_type(event: Any) → EventType[source]#

It returns the event type

get_ddl_event(event: QueryEvent) → DdlEvent[source]#

It converts the query of the BinLog into a DdlEvent. The kind of statement and the schema and table affected are taken from the query (ignoring its comments), the schema of the session is used for the tables that are not qualified. Any statement that is not a create, alter or drop of a schema or table is an DdlAction.OTHER event.

process_event(event: Any) → None[source]#

It should be implemented if another event (apart from DDL and DML) must be processed…

process_dml_event(event: Any) → list[Record][source]#

It processes the event and return the records to stream

_check_row_metadata() → None[source]#

It checks (once) if the server has binlog_row_metadata = FULL. If it does not (MINIMAL is the default of MySQL 8, and MySQL 5.7 does not have the variable) the BinLog does not have the type of the columns and mysql-replication reads the unsigned integers as negative numbers, the ENUM and SET values as None and it fails with binary or not UTF-8 data: the type of the columns is read from the database (see _apply_column_metadata).

_get_column_metadata(schema: str, table: str) → list[tuple][source]#

(data type, column type, character set) of the columns of a table, in order

_apply_column_metadata(event: Any) → None[source]#

It completes the columns of the event with what the BinLog does not have when binlog_row_metadata is not FULL, so mysql-replication decodes the rows as it does with FULL: unsigned integers, the values of ENUM and SET, binary data (bytes) and the character set of the text. The definition of the table when the rows are read is used, if the table has a different number of columns it is not applied.

static _get_json_columns(event: Any) → set[str][source]#

The names of the JSON columns of the table, as they are in the values of the rows (the position is used, UNKNOWN_COL<position>, when the name is not known).

static _decode_json_columns(values: dict, json_columns: set[str]) → dict[source]#

The values of the JSON columns have bytes for the keys and the strings, they are converted to text. A value that is bytes in any other column is binary data (e.g. a BLOB), it is not changed (Record.to_json() shows it as hexadecimal).

_get_transaction_id(begin_event: QueryEvent) → str[source]#

Identifier shared by all the rows of a transaction: <binlog file>:<position>, where the position is the one where the BEGIN event starts.

The xid (XidEvent) is not used because it is written at COMMIT, after the row events, so it can’t be attached to the rows without buffering the whole transaction; besides, it is not available for non-transactional engines (e.g. MyISAM), and it restarts with the server. The position of BEGIN is known before the first row, works with or without GTIDs and is unique within the binlog it was read from. Use global_id (GTID) for an identifier that is stable across servers.

_map_values_to_columns(values: dict, schema: str, table: str) → dict[source]#

Map UNKNOWN_COL* keys to actual column names.

Parameters:
  • values – Dictionary with potentially UNKNOWN_COL keys

  • schema – Database/schema name

  • table – Table name

Returns:

Dictionary with proper column names

_fetch_table_columns(schema: str, table: str) → list[str][source]#

Fetch column names from the database for a specific table.

Parameters:
  • schema – Database/schema name

  • table – Table name

Returns:

List of column names in order

abstractmethod _update_log_file(log_file_name: str)[source]#

Persist the current binlog file name so the stream can resume after restart.

abstractmethod _update_log_pos(position: int)[source]#

Persist the current binlog position so the stream can resume after restart.

_abc_impl = <_abc._abc_data object>#
_impls: ClassVar[dict[str, type[Self]]] = {}#

How to Use#

While using library core-cdc>=1.0.2 that uses mysql-replication>=1.0.7 the value of variable binlog_row_metadata should be FULL.

Note

MINIMAL is the default of MySQL 8, and MySQL 5.7 does not have the variable: the BinLog does not have the type of the columns. mysql-replication (tested with 1.0.15) would read the unsigned integers as negative numbers (250 in a TINYINT UNSIGNED is -6), the ENUM and SET values as None and it would lose the events with binary or not UTF-8 data (a BLOB, a latin1 text). If the processor has connection_settings it reads the type of the columns from the database (information_schema), once for each table and again after a DDL query, and the rows are decoded as with FULL (verified with MySQL 8.4, 8.0 and 5.7). It uses the current definition of the tables (if a table has a different number of columns than the event it is not applied) and it needs the privilege to read information_schema; without connection_settings those values are wrong. The processor logs, once, that it does it: FULL is recommended.

Let’s create a MySql server using Docker…

docker run \
   --name=MySqlServer \
   --env=MYSQL_ROOT_PASSWORD=mysql_password \
   --volume=/var/lib/mysql \
   -p 3306:3306 \
   --restart=no \
   --runtime=runc \
   -d mysql:5.7

Check the value in the server…

SHOW VARIABLES LIKE 'binlog_row_metadata'

Update the MySQL configuration file… This file is usually named my.cnf on Unix/Linux systems and my.ini on Windows. The location of this file can vary depending on your operating system and MySQL installation method. Common locations include /etc/mysql/my.cnf, /etc/my.cnf, or /usr/local/mysql/my.cnf.

Add or modify the binlog_row_metadata option in the [mysqld] section of the configuration file. Set it to FULL to enable full metadata logging.

[mysqld]
binlog_row_metadata = FULL

If you are using Docker based on oraclelinux-slim you can use:

docker exec -it {container-name} bash
microdnf install nano
nano /etc/my.cnf

Warning

MySQL 8.4 and later: mysql-replication (tested with 1.0.15) runs SHOW MASTER STATUS to find where to start reading whenever BinLogStreamReader is created with resume_stream=False. That statement was removed in MySQL 8.4 (replaced by SHOW BINARY LOG STATUS), so the stream fails with ProgrammingError (1064, "You have an error in your SQL syntax; ... near 'MASTER STATUS'"). The mysql:latest Docker image is affected.

Either use a MySQL server up to 8.3 (e.g. mysql:8.0), or start the stream from an explicit position with resume_stream=True, as shown below. The _update_log_file / _update_log_pos methods of the processor already persist the position, so a stored one can be used here as well.

import pymysql
from pymysqlreplication import BinLogStreamReader

with pymysql.connect(**cxn_params) as connection:
    with connection.cursor() as cursor:
        # Current position; use `SHOW MASTER STATUS` on MySQL < 8.4.
        cursor.execute("SHOW BINARY LOG STATUS")
        log_file, log_pos = cursor.fetchone()[:2]

stream = BinLogStreamReader(
    resume_stream=True,
    log_file=log_file,
    log_pos=log_pos,
    connection_settings=cxn_params,
    blocking=True,
    freeze_schema=False,
    server_id=1)

To read everything that is still in the binary logs, use the first file listed by SHOW BINARY LOGS and log_pos=4 instead.

Then, the below example script showcase how to process the MySQL BinLog…

import logging
import os
from pprint import pprint
from typing import Any, List

from core_mixins.logger import get_logger
from pymysqlreplication import BinLogStreamReader

from core_cdc.base import Record
from core_cdc.processors.mysql import MySqlBinlogProcessor
from core_cdc.targets.base import ITarget

cxn_params = {
    "host": "localhost",
    "port": 3306,
    "user": "root",
    "passwd": "mysql_password"
}

logger = get_logger(
    log_level=int(os.getenv("LOGGER_LEVEL", str(logging.INFO))),
    reset_handlers=True)


class CustomMySqlBinlogProcessor(MySqlBinlogProcessor):
    """ Custom class to implement required methods """

    def process_dml_event(self, event: Any) -> List[Record]:
        recs = super().process_dml_event(event)
        logger.info("The following records will be processed...")

        for rec in recs:
            pprint(rec.to_json())

        return recs

    def _update_log_file(self, log_file_name: str):
        # Persist it (DB, Redis, file, etc.) to resume after a restart...
        ...

    def _update_log_pos(self, position: int):
        # Persist it (DB, Redis, file, etc.) to resume after a restart...
        ...


class CustomTarget(ITarget):
    def _save(self, records: List[Record], **kwargs):
        logger.info(f"Saving: {records}")


try:
    target = CustomTarget(
        logger=logger, execute_ddl=True,
        send_data=True)

    stream = BinLogStreamReader(
        resume_stream=False,  # MySQL 8.4+: see the warning above
        connection_settings=cxn_params,
        blocking=True,
        freeze_schema=False,
        server_id=1)

    processor = CustomMySqlBinlogProcessor(
        stream=stream,
        targets=[target],
        connection_settings=cxn_params,
        service=os.getenv("SERVICE_NAME", "Functional-Tests"),
        logger=logger)

    processor.execute()

except Exception as error:
    logger.error(f"An error has been raised. Error: {error}.")

You will see something like…

[INFO] connection_settings: {'host': 'localhost', 'port': 3306, 'user': 'root', 'passwd': 'mysql_password', 'charset': 'utf8'}
[INFO] blocking: True
[INFO] allowed_events_in_packet: frozenset({<class 'pymysqlreplication.event.GtidEvent'>, <class 'pymysqlreplication.event.RandEvent'>, <class 'pymysqlreplication.event.StopEvent'>, <class 'pymysqlreplication.event.MariadbGtidListEvent'>, <class 'pymysqlreplication.event.QueryEvent'>, <class 'pymysqlreplication.row_event.TableMapEvent'>, <class 'pymysqlreplication.row_event.UpdateRowsEvent'>, <class 'pymysqlreplication.event.FormatDescriptionEvent'>, <class 'pymysqlreplication.row_event.WriteRowsEvent'>, <class 'pymysqlreplication.row_event.DeleteRowsEvent'>, <class 'pymysqlreplication.event.MariadbAnnotateRowsEvent'>, <class 'pymysqlreplication.event.ExecuteLoadQueryEvent'>, <class 'pymysqlreplication.event.MariadbStartEncryptionEvent'>, <class 'pymysqlreplication.event.HeartbeatLogEvent'>, <class 'pymysqlreplication.event.XAPrepareEvent'>, <class 'pymysqlreplication.event.MariadbGtidEvent'>, <class 'pymysqlreplication.event.MariadbBinLogCheckPointEvent'>, <class 'pymysqlreplication.event.BeginLoadQueryEvent'>, <class 'pymysqlreplication.event.UserVarEvent'>, <class 'pymysqlreplication.event.XidEvent'>, <class 'pymysqlreplication.row_event.PartialUpdateRowsEvent'>, <class 'pymysqlreplication.event.RowsQueryLogEvent'>, <class 'pymysqlreplication.event.RotateEvent'>, <class 'pymysqlreplication.event.PreviousGtidsEvent'>})
[INFO] server_id: 1
[INFO] Reading events from the stream...
[WARNING]
                    Before using MARIADB 10.5.0 and MYSQL 8.0.14 versions,
                    use python-mysql-replication version Before 1.0 version
[INFO] Received event: RotateEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 4.
[INFO] NEXT FILE: c8db74e52957-bin.000002. POSITION: 4.
[INFO] Received event: FormatDescriptionEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 123.
[INFO] Received event: PreviousGtidsEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 154.

Every record carries two identifiers: global_id is the GTID of the transaction (None when GTIDs are disabled) and transaction_id is <binlog file>:<position where BEGIN starts> (e.g. mysql-bin.000001:1560), shared by all the rows of a transaction. Use global_id when you need an identifier that is stable across servers.

The processor also converts every DDL query of the BinLog into a DdlEvent before sending it to the targets that have execute_ddl=True: the action (create, alter or drop of a schema or table, OTHER for anything else), the schema_name and table_name (taken from the query, the schema of the session is used for tables that are not qualified) and the original statement. See the Targets section to translate it to another dialect.

The queries that delimit transactions (BEGIN, COMMIT, ROLLBACK, savepoints and XA statements) are stored as queries in the BinLog but they are not schema changes, so they are EventType.GLOBAL events and never reach the targets.

Let’s execute some DDL and DML statements and follow the output in the console…

Create database

CREATE DATABASE IF NOT EXISTS test_database;
[INFO] Received event: QueryEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 418.
[INFO] The below query was executed in: CustomTarget.
[INFO] /* ApplicationName=DBeaver 24.3.0 - SQLEditor <Script-4.sql> */ CREATE DATABASE IF NOT EXISTS test_database

Create table

CREATE TABLE person (
    id INT AUTO_INCREMENT PRIMARY KEY,
    first_name VARCHAR(50) NOT NULL,
    last_name VARCHAR(50) NOT NULL,
    email VARCHAR(100),
    birth_date DATE,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
[INFO] Received event: QueryEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 872.
[INFO] The below query was executed in: CustomTarget.
[INFO] /* ApplicationName=DBeaver 24.3.0 - SQLEditor <Script-4.sql> */ CREATE TABLE person (
    id INT AUTO_INCREMENT PRIMARY KEY,
    first_name VARCHAR(50) NOT NULL,
    last_name VARCHAR(50) NOT NULL,
    email VARCHAR(100),
    birth_date DATE,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
)

Inserting records

INSERT INTO person (first_name, last_name, email, birth_date)
VALUES
('John', 'Doe', 'john.doe@example.com', '1990-01-15'),
('Jane', 'Smith', 'jane.smith@example.com', '1985-07-22');
[INFO] Received event: QueryEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 1018.
[INFO] Received event: TableMapEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 1088.
[INFO] Received event: WriteRowsEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 1211.
[INFO] The following records will be processed...
{'event_timestamp': 1733197735,
 'event_type': 'INSERT',
 'global_id': None,
 'position': 1211,
 'primary_key': '',
 'record': {'id': 1,
            'first_name': 'John',
            'last_name': 'Doe',
            'email': 'john.doe@example.com',
            'birth_date': '1990-01-15',
            'created_at': '2024-12-03T03:48:55'},
 'schema_name': 'test_database',
 'service': 'Functional-Tests',
 'source': 'c8db74e52957-bin.000002',
 'table_name': 'person',
 'transaction_id': 'c8db74e52957-bin.000002:1042'}
{'event_timestamp': 1733197735,
 'event_type': 'INSERT',
 'global_id': None,
 'position': 1211,
 'primary_key': '',
 'record': {'id': 2,
            'first_name': 'Jane',
            'last_name': 'Smith',
            'email': 'jane.smith@example.com',
            'birth_date': '1985-07-22',
            'created_at': '2024-12-03T03:48:55'},
 'schema_name': 'test_database',
 'service': 'Functional-Tests',
 'source': 'c8db74e52957-bin.000002',
 'table_name': 'person',
 'transaction_id': 'c8db74e52957-bin.000002:1042'}
[INFO] Saving: [<core_cdc.base.Record object at 0x7281f3b55340>, <core_cdc.base.Record object at 0x7281f3b55370>]
[INFO] 2 records were sent to: CustomTarget!
[INFO] {'target': 'CustomTarget', 'number_of_records': 2, 'event_type': <EventType.INSERT: 'INSERT'>, 'schema': 'test_database', 'table': 'person'}
[INFO] Received event: XidEvent.
[INFO] File: c8db74e52957-bin.000002, Position: 1242.

Updating a record

UPDATE person
SET first_name = 'Jonathan', last_name = 'Dover'
WHERE id = 1;