[ 
https://issues.apache.org/jira/browse/FLINK-40314?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Leonard Xu reassigned FLINK-40314:
----------------------------------

    Assignee: sanqingleo

> Introduce a MariaDB CDC connector
> ---------------------------------
>
>                 Key: FLINK-40314
>                 URL: https://issues.apache.org/jira/browse/FLINK-40314
>             Project: Flink
>          Issue Type: New Feature
>          Components: Flink CDC
>            Reporter: sanqingleo
>            Assignee: sanqingleo
>            Priority: Major
>              Labels: Flink-CDC
>
> h2. Motivation
> The existing MySQL CDC connector can connect to some MariaDB deployments 
> through the MySQL-compatible protocol. However, it does not correctly support 
> MariaDB GTID parsing or recovery across physical MariaDB servers.
> MySQL and MariaDB use incompatible GTID formats:
>  * MySQL: UUID-based `uuid:interval`
>  * MariaDB: `domain-server-sequence`
> When a MariaDB cluster is scaled, migrated, or failed over, or when a Flink 
> job restarts, the connector may be routed to another server in the same GTID 
> domain. The new server can have a different server ID, and its replicated 
> GTID position may temporarily lag behind the position stored in the Flink 
> checkpoint.
> The current MySQL CDC GTID implementation cannot correctly handle this 
> situation because it interprets GTIDs using the MySQL format and semantics.
> PR apache/flink-cdc#4468 has verified MariaDB GTID support against MariaDB 
> 11.4 by adding MariaDB-specific behavior directly to `mysql-cdc`. However, 
> modifying the MySQL connector introduces regression risks to one of the most 
> critical Flink CDC source connectors.
> Therefore, this issue proposes introducing a standalone 
> `flink-connector-mariadb-cdc` connector, following an implementation approach 
> similar to `oceanbase-cdc`.
> h2. Why a standalone connector
>  * Stability isolation: MariaDB-specific changes should not affect the MySQL 
> CDC execution path.
>  * Independent evolution: MySQL and MariaDB have diverged significantly, 
> especially after MySQL 8.4, and are expected to evolve in different 
> directions.
>  * Upstream precedent: Debezium introduced a standalone 
> `debezium-connector-mariadb` after extracting shared binlog functionality.
>  * Simpler configuration: A MariaDB-only connector does not require 
> `scan.dialect` or dialect-selection branches.
> h2. Proposed Design
> h3. Module structure
> Add the following modules:
>  * `flink-connector-mariadb-cdc`
>  * `flink-sql-connector-mariadb-cdc`
> The new connector should be registered under:
> `flink-cdc-connect/flink-cdc-source-connectors`
> The initial module will depend on the unmodified upstream 
> `flink-connector-mysql-cdc` and reuse its snapshot and source implementation.
> {code:java}
> flink-connector-mariadb-cdc
>   ├── flink-connector-mysql-cdc
>   ├── mysql:mysql-connector-java
>   └── flink-cdc-common
> {code}
> MariaDB behavior will be injected entirely from the new connector without 
> modifying the `mysql-cdc` source code.
> h3. Public API
>  * SQL factory identifier: `mariadb-cdc`
>  * Flink package: `org.apache.flink.cdc.connectors.mariadb`
>  * DataStream API: reuse `MySqlSource.builder()` in the first version
>  * Connector options: reuse the existing MySQL CDC options, including startup 
> mode, server ID, and incremental snapshot options
>  * Do not introduce `scan.dialect`
>  * For `scan.startup.mode=specific-offset`, the GTID must use the MariaDB 
> `domain-server-sequence` format
> A dedicated `MariaDbSource` DataStream API can be introduced in a follow-up 
> issue.
> h3. MariaDB GTID behavior
> The connector should:
>  * Enable MariaDB replication capability `4` so that `MARIADB_GTID` events 
> are emitted and the consumed GTID position advances.
>  * Read the local MariaDB GTID position from `@@gtid_binlog_pos`.
>  * Obtain the binlog file and position using `SHOW MASTER STATUS`.
>  * Parse MariaDB GTIDs using the `domain-server-sequence` format.
>  * Compare GTID positions by domain and sequence.
>  * Ignore the MariaDB server ID when evaluating GTID containment.
>  * Validate restored GTIDs and reject invalid MariaDB GTID formats.
>  * Preserve GTID information during offset serialization and checkpoint 
> recovery.
> h3. Classes to override
> The initial implementation should override the following classes from the 
> MySQL CDC dependency:
>  * `io.debezium.connector.mysql.MySqlStreamingChangeEventSource`
>  * `org.apache.flink.cdc.connectors.mysql.debezium.DebeziumUtils`
>  * 
> `org.apache.flink.cdc.connectors.mysql.debezium.task.context.StatefulTaskContext`
>  * `org.apache.flink.cdc.connectors.mysql.source.offset.BinlogOffset`
> The connector should introduce the following MariaDB-specific classes:
>  * `mariadb.debezium.MariaDBBinaryLogClient`
>  * `mariadb.debezium.MariaDbConnections`
>  * `mariadb.source.offset.MariaDbGtidStrategy`
>  * `mariadb.table.MariaDbTableSourceFactory`
> h2. Recovery Semantics
> Let:
>  * `C` be the GTID position stored in the Flink checkpoint.
>  * `H` be the executed GTID position reported by the target server through 
> `@@gtid_binlog_pos`.
>  * `L` be the earliest position still available in the target server's binlog.
> The initial implementation will only perform GTID containment validation. It 
> will not modify, downgrade, or automatically correct the checkpointed 
> position.
> If `C > H`, meaning that the selected replica temporarily lags behind the 
> checkpointed position, task initialization should fail while preserving the 
> original checkpoint position.
> The connector will not implement internal waiting, repeated position queries, 
> timeout handling, or replica reselection in the first version. Recovery will 
> rely on Flink's task restart strategy. In production, the replica normally 
> catches up within a short period, allowing a subsequent restart to recover 
> successfully.
> When reading from a replica, `log_slave_updates=ON`, or the equivalent cloud 
> database setting, is required so that replicated transactions are written to 
> the local binlog.
> h2. Known Classloader Limitation
> The first implementation will override classes using the same fully qualified 
> class names as classes contained in `mysql-cdc`.
> Therefore, `flink-sql-connector-mysql-cdc` and 
> `flink-sql-connector-mariadb-cdc` must not be placed in the same shared 
> classloader, such as `FLINK_HOME/lib`.
> If both connectors are loaded by the same classloader, JVM classpath ordering 
> determines which implementation is used:
>  * If the MySQL implementation is loaded first, `mariadb-cdc` may silently 
> lose MariaDB replication capability `4`, preventing the consumed GTID from 
> advancing.
>  * If the MariaDB implementation is loaded first, `mysql-cdc` may execute 
> MariaDB-specific GTID queries and comparison logic.
> The recommended deployment is either:
>  * Use separate Flink deployments for MySQL CDC and MariaDB CDC; or
>  * Keep both connector JARs out of `FLINK_HOME/lib` and expose only the 
> required connector through each job's user classloader.
> Supporting `mysql-cdc` and `mariadb-cdc` in the same classloader is not part 
> of the initial scope.
> h2. Testing Plan
> MariaDB 11.4 will be used as the initial integration-test baseline.
> Tests should cover:
>  * Initial snapshot reading
>  * Snapshot-to-binlog transition
>  * Savepoint and checkpoint recovery
>  * MariaDB GTID domain and sequence comparison
>  * GTID containment while ignoring server ID
>  * Binlog offset serialization and recovery
>  * Invalid GTID rejection
>  * `MARIADB_GTID` event handling
>  * Cross-server-ID recovery
> Cross-server-ID recovery can be tested deterministically without a three-node 
> cluster:
>  # Create a savepoint.
>  # Change the MariaDB server ID using `SET GLOBAL server_id`.
>  # Commit a transaction to an uncaptured table before restoring the job so 
> that `@@gtid_binlog_pos` contains the new server ID.
>  # Restore the job and verify that recovery succeeds across the server-ID 
> boundary.
> The documentation will initially declare MariaDB 10.x and 11.x support. 
> Integration coverage for MariaDB 10.6 and 10.11 LTS can be added in follow-up 
> work.
> h2. Non-Goals
>  * Do not modify the `mysql-cdc` source code.
>  * Do not package this connector in `flink-cdc-dist`.
>  * Do not switch to MariaDB Connector/J in the first version; continue using 
> MySQL Connector/J because the reused MySQL CDC path currently uses 
> `jdbc:mysql://` and `com.mysql.cj.jdbc.Driver`.
>  * Do not support `mysql-cdc` and `mariadb-cdc` in the same classloader.
> h2. Implementation Plan
>  # Create the `flink-connector-mariadb-cdc` module, SQL factory, and SPI 
> registration.
>  # Add MariaDB binlog and GTID support.
>  # Add GTID comparison, containment, serialization, and recovery tests.
>  # Create the shaded `flink-sql-connector-mariadb-cdc` artifact.
>  # Add MariaDB CDC documentation.
>  # Add MariaDB source E2E tests and CI coverage.
> h2. Related Work
>  * FLINK-34807
>  * apache/flink-cdc#4468
>  * apache/flink-cdc#2494



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to