Flink CDC 对接 PolarDB-X:用 mysql-cdc 连接器搭建增量订阅到 Elasticsearch 的实时管道

发布时间:2026/9/17 20:27:17
Flink CDC 对接 PolarDB-X:用 mysql-cdc 连接器搭建增量订阅到 Elasticsearch 的实时管道 Flink CDC 对接 PolarDB-X用 mysql-cdc 连接器搭建增量订阅到 Elasticsearch 的实时管道【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdcFlink CDC 的mysql-cdc连接器原生兼容 PolarDB-X 2.0.x 的 Binlog 协议可以像订阅 MySQL 一样实现快照 增量的全量同步。本文以官方教程「PolarDB-X 导入 Elasticsearch」为主线完整演示如何用 Docker 拉起 PolarDB-X 环境、准备订单/商品数据、在 Flink SQL CLI 中定义 CDC Source 与 ES Sink并通过 JOIN 把orders表实时丰富后写入 Elasticsearch同时结合仓库中的集成测试源码剖析 Flink CDC 为 PolarDB-X 特殊 DDL 提供的 Schema 解析兜底机制帮助你在生产环境中更稳妥地接入分布式数据库的增量变更流。为什么 PolarDB-X 可以直接用 mysql-cdc 连接器订阅PolarDB-X 兼容 MySQL 协议并暴露 Binlog 订阅通道因此无需独立连接器。在 mysql-cdc 连接器文档 的兼容性说明中明确列出了Connector支持的数据库Drivermysql-cdcMySQL 5.7 / 8.0.x / 8.4、RDS MySQL、PolarDB MySQL、Aurora MySQL、MariaDB 10.x、PolarDB X 2.0.1JDBC Driver: 8.0.27由于 PolarDB-X 的分库分表语法dbpartition by ...、tbpartition by ...、AUTO_INCREMENT BY GROUP等会让标准 MySQL DDL 解析器解析失败Flink CDC 在 Schema 构建路径中加入了兜底逻辑。从源码结构看MySqlSchema.buildTableSchema 会先用SHOW CREATE TABLE解析 DDL 构建表结构一旦解析结果中没有命中目标表PolarDB-X 的分区 DDL 常会触发这种情况自动降级执行DESC TABLE重建列结构两次都失败才会抛出FlinkRuntimeException。这一机制正是仓库中 PolarDB-X 集成测试用例要验证的核心能力见 PolardbxSourceITCase 的注释Database Polardbx supported the mysql protocol, but there are some different features in ddl. So we added fallback inMySqlSchemawhen parsing ddl failed.理解这一点后本教程的操作路径就是把 PolarDB-X 当作一个带分布式 DDL 的 MySQL用标准的mysql-cdcSource 配置即可工作。准备教程组件Docker 环境假设你在一台 MacOS 或 Linux 机器上并已安装 Docker本教程中 Elasticsearch 镜像要求 7.6.0Kibana 同为 7.6.0PolarDB-X 使用polardbx/polardb-x:2.0.1镜像与连接器文档声明的兼容版本一致。创建docker-compose.ymlversion: 2.1 services: polardbx: 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该 Docker Compose 包含三个容器PolarDB-X商品表products和订单表orders存储在该数据库中两张表将做关联得到信息更完整的订单表enriched_ordersElasticsearch最终的enriched_orders写入目标Kibana可视化查看 Elasticsearch 中的数据。在docker-compose.yml所在目录执行docker-compose up -d该命令以 detached 模式启动所有容器。可用docker ps观察容器是否正常运行访问http://localhost:5601/确认 Kibana 可用。值得对照的是仓库集成测试 PolardbxSourceTestBase 通过 Testcontainers 启动同一个polardbx/polardb-x镜像时会轮询SHOW MASTER STATUS直到 binlog 文件名非空且位点大于 4 才认为 CDC 节点就绪——如果你在本地发现 Source 启动后迟迟不吐数据可以先手动执行SHOW MASTER STATUS检查 PolarDB-X 的 binlog 是否真正可用。准备数据在 PolarDB-X 中建表并写入初始数据使用镜像内置的polardbx_root / 123456账号登录 PolarDB-X注意端口是 8527 而非 MySQL 默认的 3306mysql -h127.0.0.1 -P8527 -upolardbx_root -p123456执行以下 SQL 创建商品表与订单表并写入初始数据CREATE DATABASE mydb; USE mydb; -- 创建一张产品表并写入一些数据 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从自增 101 起、orders从 10001 起为后续增量验证订单 10004 的插入/更新/删除预留了可预期的主键。下载 Flink 与依赖 JAR 并启动集群下载 Flink 1.17.0 并解压至目录flink-1.17.0下载下列两个连接器 JAR放入flink-1.17.0/lib/目录flink-sql-connector-mysql-cdc用于订阅 PolarDB-X Binlog从 mysql-cdc 文档 可知已发布版本可在 Maven 中央仓库获取flink-sql-connector-elasticsearch7-3.0.1-1.17.jarFlink 官方的 Elasticsearch 7 Sink注意官方下载链接只对已发布版本有效SNAPSHOT 版本需要本地自行编译。启动 Flink 服务./bin/start-cluster.sh访问http://localhost:8081/可以看到 Flink 正常运行即上文配图所示的 JobManager 页面。启动 Flink SQL CLI./bin/sql-client.sh在 Flink SQL CLI 中创建表并提交实时 JOIN 作业在 SQL CLI 中依次执行以下语句。两个 Source 表都使用mysql-cdc连接器指向 PolarDB-X127.0.0.1:8527Sink 使用elasticsearch-7连接器写入enriched_orders索引整条链路不需要写一行 Java 代码-- 设置间隔时间为3秒 Flink SQL SET execution.checkpointing.interval 3s; -- 创建source1 -订单表 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 ); -- 创建source2 -产品表 Flink SQL 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 ); -- 创建sink - 关联后的结果表 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 ); -- 执行读取和写入 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;几个参数值得注意port 8527是 PolarDB-X 容器暴露的 MySQL 兼容端口与测试基类中的INNER_PORT 8527一致execution.checkpointing.interval 3s对 CDC Source 尤为关键——mysql-cdc的快照分片读取与 binlog 位点推进都依赖 checkpoint 机制间隔越小全量切增量阶段的位点越及时若 Source 并行度大于 1建议按 mysql-cdc 文档 的说明为每个 reader 设置唯一的server-id区间例如通过 SQL Hints 指定server-id 5401-5404避免多个作业/reader 共享同一 Server id 导致从错误位点读取。作业提交后全量快照先被读入随后切换到 Binlog 增量模式。在 Kibana 中查看同步结果访问http://localhost:5601/app/kibana#/management/kibana/index_pattern创建 index patternenriched_orders之后在http://localhost:5601/app/kibana#/discover即可看到orders与products关联后的 3 条订单记录。修改源表数据验证增量变更实时同步回到 PolarDB-X依次执行以下操作每执行一步刷新一次 Kibana订单数据将实时更新INSERT INTO orders VALUES (default, 2020-07-30 15:22:00, Jark, 29.71, 104, false); UPDATE orders SET order_status true WHERE order_id 10004; DELETE FROM orders WHERE order_id 10004;这三条语句恰好覆盖了Iinsert、-U/Uupdate 前后镜像、-Ddelete三类变更事件可以在 Kibana Discover 页观察到订单 10004 从出现、状态变更到最终消失的完整过程直观验证了 PolarDB-X Binlog 的增量订阅链路。源码佐证Flink CDC 如何兼容 PolarDB-X 的特殊 DDL本教程使用普通 DDL 建表但在真实 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;其中AUTO_INCREMENT BY GROUP、dbpartition by RANGE_HASH(...)、tbpartition by ...、GLOBAL INDEX均为 PolarDB-X 扩展语法。针对这类 DDLPolardbxSourceITCase 提供了系统性验证testSingleKey对单主键分区表开启scan.incremental.snapshot.enabled true做增量快照扫描先断言快照阶段读出的 5 条I记录再断言 sink 侧数据testMultiPks等用例覆盖复合主键分区表字符集用例 PolardbxCharsetITCase验证不同字符集下 binlog 数据解码的正确性。这些用例与教程中的手动验证互为印证教程证明了标准 DDL 下 PolarDB-X → ES 全链路可跑通而集成测试进一步证明了分布式分区 DDL、多主键、特殊字符集下 Schema 解析与数据读取依然正确。如果你在接入真实 PolarDB-X 集群时遇到Cant obtain schema for table ...类报错可回到 MySqlSchema 的SHOW CREATE TABLE → DESC两级解析逻辑排查目标表 DDL 是否命中了已知解析限制。环境清理在docker-compose.yml所在目录停止并删除所有容器docker-compose down进入 Flink 部署目录停止 Flink 集群./bin/stop-cluster.sh适用前提与小结本教程的适用前提PolarDB-X 2.0.1镜像polardbx/polardb-x:2.0.1、Flink 1.17.0、ES/Kibana 7.6.0订阅依赖 PolarDB-X 的 Binlog 通道可用可用SHOW MASTER STATUS检查。教程覆盖的核心能力可归纳为用标准mysql-cdcSource 配置hostname/port/username/password/database-name/table-name直接订阅 PolarDB-X无需专有连接器多 Source 实时 LEFT JOIN elasticsearch-7Sink 的纯 SQL 管道搭建方法通过 INSERT/UPDATE/DELETE 三步验证增量变更的端到端时延表现结合 MySqlSchema 兜底逻辑与 PolardbxSourceITCase 测试理解 Flink CDC 对 PolarDB-X 分区 DDL、多主键表的兼容机制为接入生产环境提供依据。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询