)
Flink CDC 实战基于 PolarDB-X CDC 构建流式 ETL 管道写入 Elasticsearch完整教程【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文以仓库 docs/content/docs/connectors/flink-sources/tutorials/polardbx-tutorial.md 为核心骨架结合 mysql-cdc.md 与 flink-connector-mysql-cdc 模块中的 PolarDB-X 集成测试源码深度展开。引言本教程要解决什么问题本教程演示如何使用 Flink CDC 为PolarDB-X快速构建流式 ETLStreaming ETL管道假设你经营一个电商业务商品products与订单orders数据存储在 PolarDB-X 中业务目标是用商品表实时补全Enrich订单数据再将补全后的订单结果实时写入 Elasticsearch 供检索与可视化。整个过程完全在 Flink SQL CLI 中完成全部使用标准 SQL 语法无需编写一行 Java/Scala 代码也不需要安装任何 IDE。读完本教程你将掌握如何用 Docker Compose 一键拉起 PolarDB-X Elasticsearch Kibana 的演示环境如何用mysql-cdcconnector 将 PolarDB-X 当作 MySQL 协议源接入 Flink CDCPolarDB-X 是 Flink CDC MySQL connector 官方支持的数据库之一详见下文「原理篇」如何用 Flink SQL 完成双表 JOIN 的实时流式补全并写入 Elasticsearch如何在 PolarDB-X 中执行 Insert / Update / Delete观察变更在 Kibana 中实时生效。一、方案概览PolarDB-X 为什么能直接用 mysql-cdcPolarDB-X原 DRDS阿里云 PolarDB 分布式版兼容 MySQL 协议与语法因此 Flink CDC 的mysql-cdcconnector 可以直接对它进行数据捕获。这一点在官方文档 mysql-cdc.md 的 Supported Databases 表中明确列出mysql-cdc支持 MySQL 5.7 / 8.0.x / 8.4、RDS MySQL、PolarDB MySQL、Aurora MySQL、MariaDB 10.x 以及PolarDB X 2.0.1使用的 JDBC 驱动为 8.0.27。仓库中还提供了针对 PolarDB-X 的专门集成测试位于 PolardbxSourceITCase.java其类注释明确写道Database Polardbx supported the mysql protocol, but there are some different features in ddl. So we added fallback in MySqlSchema when parsing ddl failed and provided these cases to test.—— 也就是说PolarDB-X 虽然走 MySQL 协议但 DDL 上存在差异如分区语法、全局二级索引Flink CDC 在 DDL 解析失败时提供了 fallback 机制来兜底。测试用例覆盖了单主键表、多主键表、全字段类型表含 JSON、空间几何类型等的快照读取与 binlog 变更捕获见 polardbx_ddl_test.sql。一句话总结本教程的技术要点PolarDB-X 通过 MySQL 协议 binlog 对外提供 CDC 能力Flink CDC 的 mysql-cdc connector 负责快照Snapshot与增量binlog两种数据的无缝衔接这正是它能在纯 SQL 场景下完成实时 ETL 的底层原因。二、环境准备用 Docker Compose 拉起全部依赖组件2.1 前置条件准备一台安装了Docker的 Linux 或 macOS 电脑。本演示所需的组件全部以容器方式运行统一由docker-compose管理。2.2 编写 docker-compose.yml创建docker-compose.yml文件内容如下version: 2.1 services: polardbx: image: polardbx/polardb-x:2.0.1 container_name: polardbx ports: - 8527:8527 elasticsearch: image: elastic/elasticsearch:7.6.0 container_name: elasticsearch 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 container_name: kibana ports: - 5601:5601 volumes: - /var/run/docker.sock:/var/run/docker.sock三个容器的分工如下容器作用PolarDB-X存放products、orders两张源表是 CDC 数据源也是订单补全 JOIN 的左侧/右侧表Elasticsearch作为数据 Sink存储补全后的订单enriched_orders索引Kibana可视化 Elasticsearch 中的数据用于验证 ETL 结果在docker-compose.yml所在目录执行以下命令启动全部容器后台模式docker-compose up -d执行docker ps检查容器是否正常运行。也可以访问 http://localhost:5601/ 确认 Kibana 是否正常启动。说明PolarDB-X 容器对外映射端口为8527这也是仓库集成测试PolardbxSourceTestBase.java中使用的INNER_PORT 8527MySQL 客户端与 Flink CDC 均通过该端口连接测试中默认用户名/密码为polardbx_root/123456与本教程保持一致。2.3 准备 Flink 与所需 JAR 包下载Flink 1.18.0bin 版scala_2.12并解压到目录flink-1.18.0将以下 JAR 包放入flink-1.18.0/lib/flink-sql-connector-mysql-cdc-3.0-SNAPSHOT.jarflink-sql-connector-elasticsearch7-3.0.1-1.17.jar⚠️版本提示下载链接仅对稳定版本stable releases可用SNAPSHOT 依赖需要基于 master 或 release 分支自行构建。flink-sql-connector-elasticsearch7为稳定发布件可从 Maven Central 仓库直接获取flink-sql-connector-mysql-cdc-3.0-SNAPSHOT属于当前开发版本需按上文说明自行构建后放入 lib 目录。另外由于 Flink CDC 项目的 MySQL Connector 涉及 GPLv2 协议兼容问题预构建连接器包中不内置 MySQL 驱动如运行环境缺失需按 mysql-cdc.md 的说明手动补充mysql:mysql-connector-java:8.0.27依赖。三、准备数据在 PolarDB-X 中建表并写入样本数据3.1 进入 PolarDB-X 数据库mysql -h127.0.0.1 -P8527 -upolardbx_root -p123456实操提示为确保 Flink DDL 中database-name mydb配置与实际一致可在连接后先执行CREATE DATABASE IF NOT EXISTS mydb; USE mydb;再建表。3.2 创建商品表与订单表并填充数据-- PolarDB-X CREATE TABLE products ( id INTEGER NOT NULL AUTO_INCREMENT PRIMARY KEY, name VARCHAR(255) NOT NULL, description VARCHAR(512) ) 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);products表9 条商品数据主键id从 101 自增orders表3 条订单数据主键order_id从 10001 自增其中product_id与products.id关联order_status表示订单是否已下单。知识延伸真实生产环境中的 PolarDB-X 表通常带分布式分区与全局索引语法。仓库测试库 polardbx_ddl_test.sql 展示了典型写法例如create table orders ( id bigint not null auto_increment by group, seller_id varchar(30) DEFAULT NULL, order_id varchar(30) DEFAULT NULL, buyer_id varchar(30) DEFAULT NULL, create_time datetime DEFAULT NULL, primary key(id), GLOBAL INDEX g_i_seller(seller_id) dbpartition by hash(seller_id) ) ENGINEInnoDB DEFAULT CHARSETutf8 dbpartition by RANGE_HASH(buyer_id, order_id, 10) tbpartition by RANGE_HASH(buyer_id, order_id, 10) tbpartitions 3;其中dbpartition by/tbpartition by是 PolarDB-X 独有的分库分表语法GLOBAL INDEX用于创建全局二级索引。Flink CDC 之所以能正确解析这类 DDL正是依赖前述MySqlSchema中的 fallback 机制。四、启动 Flink 集群与 Flink SQL CLI切换到 Flink 目录cd flink-1.18.0启动 Flink 集群./bin/start-cluster.sh随后访问 http://localhost:8081/ 确认 Flink 是否正常运行界面大致如下Flink Web UI 正常运行界面本教程中启动 Flink Standalone 集群后的状态页启动 Flink SQL CLI./bin/sql-client.sh应看到 SQL Client 的欢迎界面进入Flink SQL交互提示符。五、在 Flink SQL CLI 中使用 DDL 建表5.1 开启周期性 Checkpoint首先开启每 3 秒一次的 Checkpoint保证故障恢复与端到端一致性-- Flink SQL Flink SQL SET execution.checkpointing.interval 3s;对 CDC 场景而言Checkpoint 是「恰好一次exactly-once」语义的基石Flink CDC 会定期把 binlog 读取位点持久化到 Checkpoint 中任务重启后可从最近一次位点继续消费。5.2 创建两个 CDC 源表orders 与 products-- Flink SQL Flink SQL SET execution.checkpointing.interval 3s; -- create source table2 - orders 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 127.0.0.1, port 8527, username polardbx_root, password 123456, database-name mydb, table-name orders ); -- create source table2 - products CREATE TABLE products ( id INT, name STRING, description STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 127.0.0.1, port 8527, username polardbx_root, password 123456, database-name mydb, table-name products );WITH参数说明默认值与含义参考 mysql-cdc.md 的参数表参数必填默认值说明connector是-固定为mysql-cdchostname是-PolarDB-X 所在主机port否3306PolarDB-X 端口本演示为 8527username/password是-数据库账号演示环境为polardbx_root/123456database-name是-要监听的库名支持正则表达式匹配多库table-name是-要监听的表名支持正则表达式连接器会用\.拼接 database-name 与 table-name 组成全路径正则去匹配库表全名此外还有几个在实际生产中使用频率很高的可选参数本演示未显式配置、使用默认值即可scan.incremental.snapshot.enabled默认true增量快照机制。相比旧版快照它支持快照阶段并行读取、以 chunk 为粒度做 Checkpoint且无需全局读锁FLUSH TABLES WITH READ LOCK。推荐在使用增量快照时给server-id配置区间如5400-5408区间大小需不小于并行度。server-id数据库客户端唯一 ID用于连接 MySQL 集群并读取 binlog。默认在 5400~6400 之间随机生成建议显式指定若不同作业共享同一 server-id可能从错误的 binlog 位点读取。scan.incremental.snapshot.chunk.size默认 8096快照阶段表被切分的 chunk 行数大小测试用例中常调小到 100 以加速验证。server-time-zone数据库会话时区控制 MySQL TIMESTAMP 类型如何转为 STRING例如Asia/Shanghai未设置时使用ZoneId.systemDefault()。源码佐证在仓库集成测试 PolardbxSourceITCase.java 中PolarDB-X 源表的 DDL 即通过connector mysql-cdc声明并显式配置了scan.incremental.snapshot.enabled true、scan.incremental.snapshot.chunk.size 100、server-time-zone UTC与区间形式的server-id与本节参数表完全对应。测试同时验证了 PolarDB-X 单主键表、多主键表primary key(id, order_id)以及包含 40 字段的全类型表的快照与 DML 变更捕获。5.3 创建 Elasticsearch 结果表enriched_orders最后创建用于向 Elasticsearch 写入补全订单的结果表-- Flink SQL -- create sink table - enrich_orders 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, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector elasticsearch-7, hosts http://localhost:9200, index enriched_orders );Sink 表字段 orders 全部字段 products 的name、description两个补全字段PRIMARY KEY (order_id) NOT ENFORCED用于告知 Elasticsearch 以order_id作为文档主键实现 Upsert 语义新增写入、变更覆盖、删除移除。hosts指向本地 Elasticsearch 7.6.0 的 HTTP 端口 9200index指定目标索引名enriched_orders。注elasticsearch-7是 Flink SQL 生态的 Elasticsearch 7.x connector。仓库中还提供了一套以 YAML 配置驱动的 Pipeline 版 Elasticsearch connector见 elasticsearch.md它面向 Flink CDC 2.x 的 YAML Pipeline 编程模型本教程沿用 Flink SQL CLI elasticsearch-7connector 的实现方式。六、用 Flink SQL 完成订单补全并写入 Elasticsearch在 Flink SQL CLI 中提交如下 INSERT 语句用orders表LEFT JOIN products表关联键o.product_id p.id实时补全订单并写入 Elasticsearch-- Flink SQL Flink SQL INSERT INTO enriched_orders SELECT o.order_id, o.order_date, o.customer_name, o.price, o.product_id, o.order_status, p.name, p.description FROM orders AS o LEFT JOIN products AS p ON o.product_id p.id;LEFT JOIN的选择是有意为之即使某条订单的商品在products表中缺失订单本身也会被写入补全字段为空保证订单数据不丢。在 Flink 流式语义下这是一条双流 JOIN 的持续查询每当orders或products任一表出现变更事件Insert/Update/Delete结果都会增量重算并同步到 Elasticsearch。作业提交后补全后的订单数据即可在 Kibana 中查看。七、在 Kibana 中验证实时 ETL 结果7.1 创建 index pattern访问 http://localhost:5601/app/kibana#/management/kibana/index_pattern 创建 index patternenriched_orders在 Kibana 中创建 enriched_orders 索引模式的界面7.2 查看补全后的订单访问 http://localhost:5601/app/kibana#/discover 即可发现补全后的订单数据Kibana Discover 中展示的 enriched_orders 补全订单数据此时应能看到 3 条订单且每条订单已带上对应的商品名称与描述。7.3 在线变更验证实时更新接下来在 PolarDB-X 中依次执行增、改、删操作观察 Kibana 中enriched_orders的数据是否随每一步实时变化新增一条订单--PolarDB-X INSERT INTO orders VALUES (default, 2020-07-30 15:22:00, Jark, 29.71, 104, false);更新订单状态将 10004 号订单置为已下单--PolarDB-X UPDATE orders SET order_status true WHERE order_id 10004;删除订单--PolarDB-X DELETE FROM orders WHERE order_id 10004;每一步执行后Kibana Discover 中的数据都会在数秒内同步刷新新增订单出现、更新订单的order_status由false变为true、删除后对应文档从索引中移除。这正是 mysql-cdc 将 PolarDB-X binlog 中的 DML 事件IInsert、-U/UUpdate、-DDelete实时转为 Flink 变更流、再由elasticsearch-7sink 以主键 Upsert 落到 ES 的完整链路。源码佐证集成测试 PolardbxSourceITCase.java 中testSingleKey用例即按「先校验快照数据 → 再校验 Sink 数据 → 最后执行 INSERT/UPDATE/DELETE 校验 binlog 变更事件」三步验证其期望输出形如I[6, 1006, 1006, 1006, 2022-01-17T00:00]Insert、-D[...]Delete、先删后插形式的-U/UUpdate与本节的实时变更验证逻辑一致。八、清理环境演示结束后按以下步骤回收资源在docker-compose.yml所在目录停止全部容器docker-compose down在 Flink 目录flink-1.18.0下停止 Flink 集群./bin/stop-cluster.sh九、原理深化Flink CDC 如何支撑 PolarDB-X 场景9.1 底层工作链路从源码与测试可以还原出这条链路的四个关键环节快照读取Snapshot作业启动后mysql-cdc 通过 JDBC 按主键分 chunk 并行读取orders、products表的全量数据期间以 chunk 为粒度推进 Checkpointbinlog 监听Incremental快照完成后连接器以 MySQL 协议客户端身份分配唯一server-id订阅 PolarDB-X 的 binlog捕获持续产生的 DML 事件状态关联JOINFlink 将两个源表注册为有状态双流LEFT JOIN持续维护关联状态并增量产出结果Upsert 落库结果以order_id为主键通过elasticsearch-7connector 写入enriched_orders索引。9.2 PolarDB-X 兼容性的实现事实协议层面PolarDB-X 兼容 MySQL 协议Flink CDC 通过 JDBC驱动 8.0.27 binlog 完成捕获无需专用 connectorDDL 层面PolarDB-X 的分库分表语法dbpartition by/tbpartition by、GLOBAL INDEX等与原生 MySQL 不同Flink CDC 在MySqlSchema解析 DDL 失败时提供 fallback 兜底见 PolardbxSourceITCase.java 类注释类型层面测试用例testFullTypesDdl验证了 PolarDB-X 中 TINYINT~BIGINT、DECIMAL、DATETIME、JSON、BLOB 以及 POINT/GEOMETRY 等空间类型的读取映射生产上可直接参考 polardbx_ddl_test.sql 的类型与 Flink 侧 DDL 的对应关系。9.3 生产化建议基于文档与源码的合理推断为每个读取并行度分配独立server-id区间形式避免多作业共用导致 binlog 位点错乱大数据量表可调整scan.incremental.snapshot.chunk.size平衡快照速度与内存占用若数据库与会话时区不同显式设置server-time-zone避免 TIMESTAMP 转字符串产生时区偏差保持 Checkpoint 开启本教程 3s保障 exactly-once 与故障恢复SNAPSHOT 版 connector 需按仓库 master/release 分支自行构建生产环境优先选用匹配 Flink 版本的稳定发布件。延伸阅读MySQL CDC Connector 完整参数与使用指南mysql-cdc.mdElasticsearch Pipeline connectorYAML 配置方式elasticsearch.md更多 Flink SQL 教程MySQL→Kafka、MySQL→Doris、MySQL→StarRocks 等tutorialsPolarDB-X 集成测试源码PolardbxSourceITCase.java、PolardbxSourceTestBase.java、polardbx_ddl_test.sql【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考