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.