Sincronización de múltiples tablas MySQL a una única tabla con Apache Flink y CDC

Agregue las siguientes líneas al archivo my.cnf en la sección [mysqld]:

log-bin = mysql-bin
binlog-format = row

Instalación de Apache Flink

Desacrgue Apache Flink desde el sitio oficial:

wget https://www.apache.org/dyn/closer.lua/flink/flink-1.18.1/flink-1.18.1-bin-scala_2.12.tgz

Descomprima el archivo descargado:

tar -zxvf flink-1.18.1-bin-scala_2.12.tgz

Descarga de conectores CDC y JDBC

Descargue los conectores necesarios:

wget https://repo1.maven.org/maven2/org/apache/flink/flink-connector-jdbc/3.0.0-1.16/flink-connector-jdbc-3.0.0-1.16.jar
wget https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/2.3.0/flink-sql-connector-mysql-cdc-2.3.0.jar

Copie estos archivos JAR al directorio ./flink-1.18.0/lib/

Configuración de Flink

Modifique el archivo ./flink-1.18.0/conf/flink-conf.yaml:

rest.port: 8110
rest.address: 0.0.0.0
rest.bind-address: 0.0.0.0

Inicio del clúster de Flink

Ejecute los siguientes comandos:

./start-cluster.sh  # Para iniciar el clúster
./stop-cluster.sh   # Para detener el clúster
jps                 # Para verificar los procesos en ejecución

Debería ver procesos similares a:

10160 StandaloneSessionClusterEntrypoint
10648 TaskManagerRunner
22557 Jps
13951 TaskManagerRunner

Script SQL para sincronización de datos

Estructura de las tablas origen

Base de datos: wordpress, Tablas: productos (información de productos) y info_almacenamiento (información de stock)

CREATE TABLE `productos` (
  `id` int(10) unsigned NOT NULL AUTO_INCREMENT,
  `nombre` varchar(20) DEFAULT '',
  `descripcion` varchar(50) DEFAULT '',
  `fecha_creacion` datetime DEFAULT NULL,
  PRIMARY KEY (`id`)
) ENGINE=InnoDB AUTO_INCREMENT=7 DEFAULT CHARSET=utf8mb4;

CREATE TABLE `info_almacenamiento` (
  `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
  `id_producto` bigint(20) unsigned NOT NULL DEFAULT '0',
  `cantidad` int(20) unsigned NOT NULL DEFAULT '0',
  PRIMARY KEY (`id`)
) ENGINE=InnoDB AUTO_INCREMENT=6 DEFAULT CHARSET=utf8mb4;

Estructura de la tabla destino

Base de datos: wordpress_destino, Tabla: producto_stock

CREATE TABLE `producto_stock` (
  `id` bigint(20) unsigned NOT NULL AUTO_INCREMENT,
  `id_producto` bigint(20) unsigned NOT NULL,
  `nombre_producto` varchar(20) NOT NULL DEFAULT '',
  `cantidad_disponible` int(10) unsigned NOT NULL DEFAULT '0',
  `fecha_creacion` datetime NOT NULL,
  UNIQUE KEY `uniq_id_producto` (`id_producto`),
  KEY `id` (`id`)
) ENGINE=InnoDB AUTO_INCREMENT=6 DEFAULT CHARSET=utf8mb4;

Script Flink SQL

Cree un archivo llamado sincronizacion_mysql.sql en el directorio /bin/:

SET execution.checkpointing.interval = 60s;

DROP TABLE IF EXISTS productos_origen;
CREATE TABLE productos_origen (
  id INT NOT NULL,
  nombre STRING,
  descripcion STRING,
  fecha_creacion TIMESTAMP,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'servidor-mysql.ejemplo.com',
  'port' = '3306',
  'username' = 'usuario_flk',
  'password' = 'contraseña_segura',
  'database-name' = 'wordpress',
  'server-time-zone' = 'America/Mexico_City',
  'table-name' = 'productos'
);

DROP TABLE IF EXISTS almacenamiento_origen;
CREATE TABLE almacenamiento_origen (
  id INT NOT NULL,
  id_producto INT,
  cantidad INT,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'servidor-mysql.ejemplo.com',
  'port' = '3306',
  'username' = 'usuario_flk',
  'password' = 'contraseña_segura',
  'database-name' = 'wordpress',
  'server-time-zone' = 'America/Mexico_City',
  'table-name' = 'info_almacenamiento'
);

DROP TABLE IF EXISTS producto_stock_destino;
CREATE TABLE producto_stock_destino (
  id_producto INT,
  nombre_producto STRING,
  cantidad_disponible INT,
  fecha_creacion TIMESTAMP,
  PRIMARY KEY (id_producto) NOT ENFORCED
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://servidor-mysql.ejemplo.com:3306/wordpress_destino?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=UTC',
  'username' = 'usuario_flk',
  'password' = 'contraseña_segura',
  'table-name' = 'producto_stock',
  'driver' = 'com.mysql.cj.jdbc.Driver',
  'scan.fetch-size' = '200'
);

INSERT INTO producto_stock_destino
SELECT p.id AS id_producto, 
       p.nombre AS nombre_producto, 
       a.cantidad AS cantidad_disponible, 
       p.fecha_creacion 
FROM productos_origen AS p 
JOIN almacenamiento_origen AS a ON p.id = a.id_producto;

Ejecución del script

Ejecute el siguiente comando para procesar el script SQL:

./bin/sql-client.sh -f ./bin/sincronizacion_mysql.sql

Verifique los datos en la base de datos destino y observe cómo se sincronizan los cambios en tiempo real.

Etiquetas: Apache Flink Flink CDC MySQL Sincronización de datos Procesamiento de flujos

Publicado el 10-8 22:04