
用 Flink CDC 构建 MySQL PostgreSQL 到 Elasticsearch 的实时流式 ETL【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本教程源自 Flink CDC 官方文档《Building a Streaming ETL with Flink CDC》build-streaming-etl-tutorial.md演示如何基于 Flink CDC 的 MySQL Connector 与 Postgres Connector以纯 SQL 方式快速搭建一套跨库实时 ETL将 MySQL 中的商品、订单数据与 PostgreSQL 中的物流数据实时关联Streaming Join增强后的订单实时写入 Elasticsearch 并通过 Kibana 可视化。整个过程全部在 Flink SQL CLI 中完成不需要编写任何一行 Java/Scala 代码也不需要安装 IDE。读完本文你将掌握用docker-compose一键拉起 MySQL、PostgreSQL、Elasticsearch 与 Kibana 测试环境在 Flink SQL CLI 中通过 DDL 声明mysql-cdc、postgres-cdc源表和elasticsearch-7结果表用一条INSERT INTO ... SELECT ... LEFT JOIN实现订单的实时流式增强并理解增量快照、binlog/复制槽消费等底层原理。一、业务场景与整体架构假设我们运营一个电商业务MySQL中存放商品表products与订单表ordersPostgreSQL中存放与订单关联的物流表shipments我们期望使用products与shipments对订单进行实时增强补全商品名称、描述以及物流的发货地、目的地、是否送达等字段再把增强后的订单实时加载到Elasticsearch中供检索与可视化。整体数据流为经典的流式 ETL 三阶段Extract抽取mysql-cdc与postgres-cdc分别订阅 MySQL binlog 与 PostgreSQL WALWrite-Ahead Log/复制槽实时捕获源库的数据变更Transform转换Flink CDC 解析变更事件通过 Streaming Join 完成跨库多表关联Load加载将增强后的订单以 Upsert 语义写入 Elasticsearch 索引enriched_orders由 Kibana 实时展示。二、准备环境Docker Compose 一键启动所需组件本教程全部组件以容器方式管理准备一台安装了 Docker 的 Linux 或 macOS 机器创建docker-compose.yml文件并填入以下内容version: 2.1 services: postgres: image: debezium/example-postgres:1.1 ports: - 5432:5432 environment: - POSTGRES_DBpostgres - POSTGRES_USERpostgres - POSTGRES_PASSWORDpostgres mysql: image: debezium/example-mysql:1.1 ports: - 3306:3306 environment: - MYSQL_ROOT_PASSWORD123456 - MYSQL_USERmysqluser - MYSQL_PASSWORDmysqlpw elasticsearch: image: elastic/elasticsearch:7.6.0 environment: - cluster.namedocker-cluster - bootstrap.memory_locktrue - ES_JAVA_OPTS-Xms512m -Xmx512m - discovery.typesingle-node ports: - 9200:9200 - 9300:9300 ulimits: memlock: soft: -1 hard: -1 nofile: soft: 65536 hard: 65536 kibana: image: elastic/kibana:7.6.0 ports: - 5601:5601这套 Compose 环境由 4 个容器组成各自职责如下容器角色MySQL存放products、orders表作为 CDC 源表与 Postgres 数据关联以增强订单Postgres存放shipments表作为 CDC 源表Elasticsearch作为数据 Sink存储增强后的订单Kibana可视化 Elasticsearch 中的数据在包含docker-compose.yml的目录下执行以下命令启动全部容器detached 模式docker-compose up -d启动完成后可用docker ps检查各容器是否正常运行也可以访问 http://localhost:5601/ 确认 Kibana 是否正常。三、准备 Flink 与所需 JAR 包下载 Flink 1.18.0 并解压到目录flink-1.18.0将以下 JAR 包放入flink-1.18.0/lib/目录下下载链接仅对稳定版本有效SNAPSHOT 依赖需要基于 master 或 release 分支自行构建。flink-sql-connector-elasticsearch7-3.0.1-1.17.jarElasticsearch 7 的 SQL Connectorflink-sql-connector-mysql-cdc-3.0-SNAPSHOT.jarMySQL CDC SQL Connector对应仓库 flink-sql-connector-mysql-cdcflink-sql-connector-postgres-cdc-3.0-SNAPSHOT.jarPostgres CDC SQL Connector对应仓库 flink-sql-connector-postgres-cdc需要说明的是由于 MySQL Connector 的 GPLv2 许可与 Flink CDC 项目不兼容预构建的连接器包中不包含 MySQL 驱动。若运行环境缺失需要自行在lib/下补充mysql:mysql-connector-java:8.0.27依赖详见 mysql-cdc.md 的 Dependencies 章节。四、准备源库数据4.1 准备 MySQL 数据进入 MySQL 容器docker-compose exec mysql mysql -uroot -p123456创建库表并写入初始数据-- MySQL CREATE DATABASE mydb; USE mydb; CREATE TABLE products ( id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY, name VARCHAR(255) NOT NULL, description VARCHAR(512) ); ALTER TABLE products AUTO_INCREMENT 101; INSERT INTO products VALUES (default,scooter,Small 2-wheel scooter), (default,car battery,12V car battery), (default,12-pack drill bits,12-pack of drill bits with sizes ranging from #40 to #3), (default,hammer,12oz carpenters hammer), (default,hammer,14oz carpenters hammer), (default,hammer,16oz carpenters hammer), (default,rocks,box of assorted rocks), (default,jacket,water resistent black wind breaker), (default,spare tire,24 inch spare tire); CREATE TABLE orders ( order_id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY, order_date DATETIME NOT NULL, customer_name VARCHAR(255) NOT NULL, price DECIMAL(10, 5) NOT NULL, product_id INTEGER NOT NULL, order_status BOOLEAN NOT NULL -- Whether order has been placed ) AUTO_INCREMENT 10001; INSERT INTO orders VALUES (default, 2020-07-30 10:08:22, Jark, 50.50, 102, false), (default, 2020-07-30 10:11:09, Sally, 15.00, 105, false), (default, 2020-07-30 12:00:30, Edward, 25.25, 106, false);4.2 准备 PostgreSQL 数据进入 Postgres 容器docker-compose exec postgres psql -h localhost -U postgres创建表并写入数据-- PG CREATE TABLE shipments ( shipment_id SERIAL NOT NULL PRIMARY KEY, order_id SERIAL NOT NULL, origin VARCHAR(255) NOT NULL, destination VARCHAR(255) NOT NULL, is_arrived BOOLEAN NOT NULL ); ALTER SEQUENCE public.shipments_shipment_id_seq RESTART WITH 1001; ALTER TABLE public.shipments REPLICA IDENTITY FULL; INSERT INTO shipments VALUES (default,10001,Beijing,Shanghai,false), (default,10002,Hangzhou,Shanghai,false), (default,10003,Shanghai,Hangzhou,false);这里有一个容易被忽略但很关键的操作ALTER TABLE public.shipments REPLICA IDENTITY FULL;。PostgreSQL 逻辑解码在默认REPLICA IDENTITY DEFAULT下更新事件只携带主键列设为FULL后变更日志才会携带整行的旧值/新值从而保证下游能看到被修改的origin、destination、is_arrived等完整字段。这也解释了为什么本教程的shipments表在 Flink DDL 中声明了主键后仍能正确捕获UPDATE事件。五、启动 Flink 集群与 Flink SQL CLIcd flink-1.18.0启动 Flink 集群./bin/start-cluster.sh访问 http://localhost:8081/ 可查看 Flink Web UI确认集群正常运行界面截图参见 flink-ui.png。启动 Flink SQL CLI./bin/sql-client.sh进入后即可看到 CLI 欢迎界面截图参见 flink-sql-client.png。六、用 Flink DDL 创建 CDC 源表与 Elasticsearch 结果表6.1 开启 Checkpoint在 Flink SQL CLI 中先开启 Checkpoint每 3 秒一次。Checkpoint 是 CDC 任务提供故障恢复与精确一次语义的基础-- Flink SQL Flink SQL SET execution.checkpointing.interval 3s;6.2 创建 MySQL CDC 源表创建捕获 MySQLproducts、orders变更的表-- Flink SQL Flink SQL CREATE TABLE products ( id INT, name STRING, description STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username root, password 123456, database-name mydb, table-name products ); Flink SQL CREATE TABLE orders ( order_id INT, order_date TIMESTAMP(0), customer_name STRING, price DECIMAL(10, 5), product_id INT, order_status BOOLEAN, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username root, password 123456, database-name mydb, table-name orders );关于mysql-cdc连接器的关键参数完整参数表见 mysql-cdc.md 的 Connector Options 章节参数必填默认值说明connector是无固定为mysql-cdchostname是无MySQL 服务器地址或主机名username/password是无连接 MySQL 使用的账号密码database-name是无监听的数据库名支持正则表达式table-name是无监听的表名支持正则表达式port否3306MySQL 端口server-id否随机生成5400–6400客户端唯一 ID 或 ID 区间如5400-5408不同任务必须使用不同 server-id否则可能从错误的 binlog 位置读取scan.incremental.snapshot.enabled否true增量快照机制快照阶段可并行、按 chunk 粒度做 Checkpoint、无需FLUSH TABLES WITH READ LOCK全局读锁scan.incremental.snapshot.chunk.size否8096快照分块的行数大小scan.startup.mode否initial可选initial/earliest-offset/latest-offset/specific-offset/timestamp/snapshot需要强调server-id的语义MySQL 会为每个读取 binlog 的客户端维护一个唯一的 server id 来维持连接并记录 binlog 位置因此不同作业若共享同一个 server id可能导致从错误的 binlog 位置开始读取。建议通过 SQL Hints 为每个 reader 分配不同 ID例如源并行度为 4 时使用SELECT * FROM source_table /* OPTIONS(server-id5401-5404) */;。当启用增量快照默认开启且需要并行读取快照时server-id必须是大于并行度的区间形式如5400-6400。6.3 创建 Postgres CDC 源表创建捕获 PostgreSQLshipments变更的表-- Flink SQL Flink SQL CREATE TABLE shipments ( shipment_id INT, order_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, PRIMARY KEY (shipment_id) NOT ENFORCED ) WITH ( connector postgres-cdc, hostname localhost, port 5432, username postgres, password postgres, database-name postgres, schema-name public, table-name shipments, slot.name flink );关于postgres-cdc连接器的关键参数完整参数表见 postgres-cdc.md 的 Connector Options 章节参数必填默认值说明connector是无固定为postgres-cdchostname是无PostgreSQL 服务器地址username/password是无连接 PostgreSQL 的账号密码database-name是无监听的数据库名schema-name是无监听的 schema 名支持正则表达式table-name是无监听的表名支持正则表达式port否5432PostgreSQL 端口slot.name是无逻辑解码复制槽名称名称只能包含小写字母、数字与下划线decoding.plugin.name否decoderbufs逻辑解码插件支持decoderbufs、wal2json、wal2json_rds、wal2json_streaming、wal2json_rds_streaming、pgoutputchangelog-mode否allallretract 流或upsert表有主键时可用无需REPLICA IDENTITY FULLscan.incremental.snapshot.enabled否false增量快照实验特性开启后可并行快照、按 chunk 做 Checkpointheartbeat.interval.ms否30s心跳事件间隔用于跟踪最新可用的复制槽偏移需要特别注意的是slot.name建议为不同表设置不同的 slot 名以避免PSQLException: ERROR: replication slot flink is active for PID 974之类的冲突错误。PostgreSQL 的复制槽会持有 WAL 数据消费进度LSN只有在提交后才触发日志清理因此连接器默认会延迟 3 个 Checkpointscan.lsn-commit.checkpoints-num-delay默认 3才滚动提交 LSN 偏移以保证故障恢复时仍可访问更早的偏移。6.4 创建 Elasticsearch 结果表创建enriched_orders表用于将增强后的订单写入 Elasticsearch-- Flink SQL Flink SQL CREATE TABLE enriched_orders ( order_id INT, order_date TIMESTAMP(0), customer_name STRING, price DECIMAL(10, 5), product_id INT, order_status BOOLEAN, product_name STRING, product_description STRING, shipment_id INT, origin STRING, destination STRING, is_arrived BOOLEAN, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector elasticsearch-7, hosts http://localhost:9200, index enriched_orders );注意PRIMARY KEY (order_id) NOT ENFORCED的声明Elasticsearch Sink 会依据主键把 CDC 的增删改事件转换为对索引文档的 Upsert 操作从而保证源库的UPDATE/DELETE能正确反映到 ES 文档上。七、Streaming Join增强订单并实时写入 Elasticsearch一切就绪后执行一条纯 SQL 语句即可完成整个流式 ETL 的核心逻辑将orders与products、shipments做左连接补全商品与物流信息写入enriched_orders-- Flink SQL Flink SQL INSERT INTO enriched_orders SELECT o.*, p.name, p.description, s.shipment_id, s.origin, s.destination, s.is_arrived FROM orders AS o LEFT JOIN products AS p ON o.product_id p.id LEFT JOIN shipments AS s ON o.order_id s.order_id;提交后该作业会持续运行先对三张源表做一次一致性快照initial 模式默认行为随后持续消费 binlog 与 WAL把之后发生的每一条变更实时推进到 Elasticsearch。7.1 在 Kibana 中验证增强结果首先在 Kibana 中创建索引模式访问 http://localhost:5601/app/kibana#/management/kibana/index_pattern创建enriched_orders索引模式界面参见 kibana-create-index-pattern.png。然后访问 http://localhost:5601/app/kibana#/discover 查看增强后的订单可以看到每条订单已经携带了来自 MySQL 的商品名称、描述以及来自 PostgreSQL 的物流发货地、目的地、是否送达等字段——这正是流式 ETL 跨库 Join的成果。7.2 实时联动验证在源库中制造变更接下来在数据库中依次执行以下操作观察 Kibana 中的增强订单如何实时联动更新1. 在 MySQL 中插入一条新订单--MySQL INSERT INTO orders VALUES (default, 2020-07-30 15:22:00, Jark, 29.71, 104, false);2. 在 PostgreSQL 中插入一条对应物流记录--PG INSERT INTO shipments VALUES (default,10004,Shanghai,Beijing,false);3. 在 MySQL 中更新订单状态--MySQL UPDATE orders SET order_status true WHERE order_id 10004;4. 在 PostgreSQL 中更新物流状态--PG UPDATE shipments SET is_arrived true WHERE shipment_id 1004;5. 在 MySQL 中删除该订单--MySQL DELETE FROM orders WHERE order_id 10004;每一步操作完成后Kibana 中的enriched_orders都会在秒级内自动刷新完整联动效果参见 kibana-detailed-orders-changes.gif新订单带物流信息出现、订单状态与送达状态实时翻转、删除操作联动移除文档。这验证了 Flink CDC 不仅能同步数据还能把跨数据库的变更事件以流式 Join 的方式实时透传。八、清理环境教程完成后在docker-compose.yml所在目录停止全部容器docker-compose down在 Flink 目录flink-1.18.0下停止 Flink 集群./bin/stop-cluster.sh九、底层原理Flink CDC 是如何做到实时 精确一次的9.1 MySQL 侧增量快照Incremental Snapshot与 binlog 消费本教程中 MySQL 源表默认开启的增量快照是 Flink CDC 3.x 的核心机制。从源码结构看MySqlSource 是基于 FLIP-27 的新版 Source 接口实现支持快照阶段并行读取表快照按 chunk key 切分为多个 snapshot chunk由 MySqlSnapshotSplitAssigner 分发给多个 reader 并行读取chunk 粒度 Checkpoint快照阶段可按分块提交 Checkpoint解决旧机制下大表快照 Checkpoint 超时的问题无全局读锁采用 Watermark Signal Algorithm受 DBLog 论文启发快照前无需执行FLUSH TABLES WITH READ LOCK因此也不需要RELOAD权限。具体到单个 chunk 的读取其流程是记录当前 binlog 位置为LOW偏移 → 用SELECT ... WHERE chunk_key low AND chunk_key high读取并缓冲 chunk 记录 → 记录当前位置为HIGH偏移 → 读取LOW到HIGH之间属于该 chunk 的 binlog 记录 → 将 binlog 变更 Upsert 进缓冲并输出保证最终一致→ 之后由单一 binlog reader 继续消费该 chunk 后续变更。当所有 snapshot chunk 完成后Source 会转为单并行度消费 binlog进入 binlog 阶段MySqlBinlogSplitAssigner 负责此阶段分片管理。为了保证快照记录与 binlog 记录的全局顺序binlog reader 会等待快照完成后出现一个完整的 Checkpoint 才开始消费此后按行级状态记录已消费的 binlog 位置。配合周期 Checkpoint作业故障重启后可精确恢复到上次成功 Checkpoint 的位置实现**精确一次Exactly-Once**语义。9.2 PostgreSQL 侧逻辑复制槽与 WAL 消费与 MySQL 直接读 binlog 不同Postgres CDC 基于 PostgreSQL 的逻辑复制Logical Replication连接器通过slot.name指定的复制槽订阅 WAL 中的变更并依赖decoding.plugin.name指定的插件本教程默认decoderbufs解析变更内容。表上的REPLICA IDENTITY FULL保证更新事件携带完整旧值从而让changelog-mode all默认能输出标准I/-U/U/-D四种行类型。从 PostgresSourceBuilder 的源码结构看Postgres 连接器同样提供增量快照模式scan.incremental.snapshot.enabledtrue默认关闭、属实验特性开启后可与 MySQL 侧一样支持并行快照与 chunk 级 Checkpoint。需要留意的是当增量快照关闭时快照扫描阶段没有可恢复的位置无法执行 Checkpoint超时的 Checkpoint 会被判定失败并可能触发作业重启。若源表数据量较大建议按 postgres-cdc.md 中的说明放宽 Checkpoint 配置如将execution.checkpointing.interval调大、允许一定数量的失败 Checkpoint。9.3 流式 Join 为什么能持续产出INSERT INTO enriched_orders SELECT ... LEFT JOIN ...之所以是一条持续运行的流作业是因为三个 CDC 源表都携带主键NOT ENFORCEDFlink 会将其识别为含主键的 changelog 流从而可以在状态中维护当前最新值以lookup-free 的流式双流 Join方式对订单、商品、物流进行持续关联。任何一侧源表的 INSERT/UPDATE/DELETE 都会驱动 Join 结果更新并由 Elasticsearch Sink 按主键 upsert 到索引文档——这正是第 7.2 节五个变更步骤能实时联动反映到 Kibana 的根本原因。十、延伸阅读MySQL CDC 连接器完整文档连接器全部参数、增量快照算法、启动位置模式、GTID 高可用、无主键表支持等Postgres CDC 连接器完整文档复制槽、解码插件、增量快照实验、分区表支持等更多源连接器教程索引见 flink-sources/tutorials 目录本文教程源码路径build-streaming-etl-tutorial.md中文版见 content.zh 对应文件【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考