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:
IProcessorIt 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_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 theBEGINevent 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 ofBEGINis 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>#
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;