恒美微站 Logo 恒美微站
  • 首页
  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心
  • 联系我们

RisingWave Kafka CDC Sink 集成测试实战:Debezium 格式下 Flink SQL 与 JDBC 双管道落库方案

  • 首页
  • 资讯中心
  • /
  • RisingWave Kafka CDC Sink 集成测试实战:Debezium 格式下 Flink SQL 与 JDBC 双管道落库方案

相关资讯

IronClaw OOBE 首次运行引导:从空白落地页到 Agent 驱动的建议卡片全链路解析 2026/9/25 2:29:35
PyFlink StreamExecutionEnvironment 完全指南:从环境创建到作业提交的核心 API 实战 2026/9/25 2:29:35
003 缺失价格调研——为 OpenCodex 补齐 jawcode 目录外的官方单价 2026/9/25 2:24:35

最新资讯

Aster:Windows原生Session级副屏实现一机双桌面
iOS高版本备份降级恢复原理与实操指南
Keil uVision2安装使用教程:51单片机C51开发环境搭建避坑指南
芯片烧录固件版本管理:命名、哈希与工具链匹配避坑指南
移动机顶盒CM211-1刷机全攻略:短接、固件选择与救砖实战
泥人网络继电器TCP Server配置与AT指令实战指南

今日推荐

AI元人文:从工具使用到思维重构的深度探索
Python+CNN车牌识别实战:从数据预处理到模型训练与部署
Vim基础操作全攻略:保存退出、模式切换与高频命令实战

本周热门

BrewUI:给Homebrew套上图形界面,让macOS软件包管理更简单
BrewUI:让Homebrew包管理变得可视化与高效
公式与文本对齐全攻略:从Word到LaTeX的实用技巧

本月精选

自研推理加速器Redwood:两周内实现PyTorch模型高效部署的实战教程
V4L2摄像头采集实战:从camera_client.rar到出图全流程解析
从“谁发明了钢琴键”到知识问答智能体:RAG与记忆工程实践

RisingWave Kafka CDC Sink 集成测试实战:Debezium 格式下 Flink SQL 与 JDBC 双管道落库方案

发布时间:2026/9/25 2:29:35
RisingWave Kafka CDC Sink 集成测试实战:Debezium 格式下 Flink SQL 与 JDBC 双管道落库方案 数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载本文以 RisingWave 仓库内的 kafka-cdc-sink 集成测试 为蓝本讲解如何利用cargo make setup一键搭建一条「RisingWave → KafkaDebezium JSON→ 下游」的 CDC 落库链路一条由 Flink SQL 消费并写入 PostgreSQL另一条由 Debezium JDBC Sink Connector 写入 MySQL 与 PostgreSQL。读完本文你将掌握 RisingWave Kafka Sink 的FORMAT DEBEZIUM用法、Debezium Connect 的注册参数、Flink 端debezium-json消费配置以及一套完整的类型兼容性回归验证方法。一、测试场景与整体架构该集成测试的核心目的是验证 RisingWave 通过 Kafka 以Debezium JSON格式对外输出变更数据CDC后异构下游系统能否正确消费。场景中的两条管道都遵循相同的数据流向Risingwave -(Debezium Json)- Kafka -(Debezium Json)- Flink SQL -- Postgresql Risingwave -(Debezium Json)- Kafka -(Debezium Json)- JDBC connector - MySQL and Postgresql管道一Flink SQL 消费Flink SQL Client 以debezium-json格式读取 Kafka topiccounts通过 JDBC 连接器写入 PostgreSQL 表flinkcounts。管道二JDBC Connector 消费Debezium Connectdebezium/connect:2.4.0.Final注册JdbcSinkConnector同时将 topiccounts与types的数据分别写入 MySQL 与 PostgreSQL。值得说明的是两条管道是并列独立的上游 RisingWave 产生的 Debezium JSON 事件会被不同消费端各自以不同方式解释和落库这恰好验证了 Debezium 消息协议的通用性与互操作性。二、环境准备在运行测试前需要满足两项基础条件见 README安装cargo make—— 用于驱动整个测试任务链安装 Docker —— 用于拉起 RisingWave、Kafka、MySQL、PostgreSQL、Flink、Debezium Connect 等容器。测试目录下 Makefile.toml 定义了全部自动化任务。任务依赖关系如下任务作用依赖clean-alldocker compose down --remove-orphans -v清理旧环境—build-imagedocker compose build构建 Flink 与 Connect 镜像clean-allstart-clusterdocker compose up -d启动全部服务并等待 10 秒build-imagesetup-prepare执行 prepare.sh 创建 Kafka topic、初始化 Flink 作业与 PG Connectorstart-clustersetup-risingwave通过 psql 依次执行create_source.sql/create_mv.sql/create_sink.sqlsetup-preparesetup-connect通过 curl 注册 MySQL JDBC Sink Connectorsetup-risingwavesetup一键串起上述全部任务上述所有任务check-pg/check-mysql/check-flink-pg分别查询 PostgreSQLcounts、MySQLcounts、PostgreSQLflinkcounts—check-db依次执行上述三个检查—check-com-pg/check-flink-com-pg查询 PostgreSQLtypes/flink_types类型兼容性—容器编排docker-compose.yml 中RisingWave 单机节点、Kafkamessage_queue、Grafana、MinIO、Prometheus 均通过extends复用仓库根目录 docker/docker-compose.yml 中定义的服务其余为测试专属服务postgresPostgreSQL 镜像账号myuser/123456库名mydb对外端口5432并以-c wal_levellogical启动开启逻辑复制能力prepare_postgres一次性容器将 postgres_prepare.sql 灌入 PostgreSQL预建目标表mysqlMySQL 镜像root 密码mysql账号myuser/123456库名mydb对外端口3306flink-jobmanager / flink-taskmanager / flink-sql-client基于 flink/Dockerfile 构建的 Flink 集群JobManager Web UI 映射到宿主机8082避开被 Kafka 占用的8081并内置了flink-sql-connector-kafka-1.17.0.jar、flink-connector-jdbc-3.0.0-1.16.jar与postgresql-42.6.0.jarconnectDebezium Connect 镜像暴露8083REST API与5005BOOTSTRAP_SERVERSmessage_queue:29092指向容器网络内的 Kafka。预置目标表结构postgres_prepare.sql 定义了 4 张目标表counts(id INT, sum INT, primary key(id))—— JDBC Connector 写入的主链路目标表flinkcounts(id BIGINT, sum BIGINT, primary key(id))—— Flink SQL 写入的目标表types与flink_types—— 全字段类型表供类型兼容性回归使用types中decimal与interval以text承载以规避精度/语法差异。三、一键搭建测试管道cargo make setup在项目根目录执行cargo make setup该命令按依赖顺序执行清空旧容器 → 构建镜像 → 启动集群 → 准备 Kafka topic 与 Flink 作业 → 在 RisingWave 中建源、建物化视图、建 Sink → 注册 MySQL Connector。1. 启动阶段start-cluster→setup-prepare集群启动后prepare.sh 依次完成使用 Redpanda 的rpk工具创建三个 topiccounts2 个分区、types1 个分区、flinktypes1 个分区通过docker compose run flink-sql-client执行 flink.sql 与 compatibility-flink.sql在 Flink 侧建立读取 Kafka、写入 PG 的作业通过 curl 向localhost:8083/connectors/注册 pg-sink.jsonPostgreSQL JDBC Sink执行 compatibility-rw.sql 与compatibility-rw-flink.sql向 RisingWave 注入全类型极值数据开启类型兼容性链路。2. RisingWave 侧 DDLsetup-risingwave该任务通过 psql 依次执行三个 SQL 文件连接参数均为-h localhost -p 4566 -d dev -U root① 建源create_source.sql使用内置datagen连接器模拟业务数据源每秒生成 10 行(id, num)随机数据CREATE TABLE metrics ( id int, num int, ) WITH ( connector datagen, fields.id.kind random, fields.id.min 1, fields.id.max 10000000, fields.num.kind random, fields.num.min -100, fields.num.max 100000, datagen.rows.per.second 10 ) FORMAT PLAIN ENCODE JSON;② 建物化视图create_mv.sql实时聚合将流式数据按id分组求和CREATE MATERIALIZED VIEW counts as select id, sum(num) from metrics group by id;③ 建 Sinkcreate_sink.sqlset sink_decouple false; CREATE SINK IF NOT EXISTS counts_sink FROM counts WITH ( connector kafka, properties.bootstrap.servermessage_queue:29092, topic counts, primary_key id ) FORMAT DEBEZIUM ENCODE JSON;这里有两个关键点primary_key idSink 以id作为主键Debezium 格式输出时会携带before/after结构下游可据此做 upsert 与删除FORMAT DEBEZIUM ENCODE JSON这是整条管道的“语言契约”——上游按 Debezium 协议编码为 JSON下游 Flink 用debezium-json解析、Debezium Connect 原生兼容sink_decouple false关闭 sink 解耦使 sink 与物化视图在同一事务内推进保证测试中数据链路的强一致性语义本测试为验证一致性而设。④ 注册 MySQL Connectorsetup-connect向 Debezium Connect 提交 mysql-sink.jsoncurl -i -X POST -H Accept:application/json -H Content-Type:application/json \ http://localhost:8083/connectors/ -d jsons/mysql-sink.json四、Debezium JDBC Sink Connector 配置详解管道二依赖 Debezium 生态的io.debezium.connector.jdbc.JdbcSinkConnector。MySQL 与 PG 两个 Connector 的差异仅在于目标库与 topic 集合MySQLmysql-sink.json{ name: mysql-jdbc-sink, config: { connector.class: io.debezium.connector.jdbc.JdbcSinkConnector, tasks.max: 1, topics: counts, connection.url: jdbc:mysql://mysql:3306/mydb, connection.username: myuser, connection.password: 123456, transforms: unwrap, transforms.unwrap.type: io.debezium.transforms.ExtractNewRecordState, transforms.unwrap.drop.tombstones: false, auto.create: true, insert.mode: upsert, delete.enabled: true, schema.evolution: basic, primary.key.fields: id, primary.key.mode: record_key } }PostgreSQLpg-sink.json结构与上几乎一致但topics为counts, types同时消费主链路与类型兼容性链路且auto.create设为false—— 因为 PG 侧目标表已由postgres_prepare.sql预建。关键参数语义参数值说明connector.classio.debezium.connector.jdbc.JdbcSinkConnectorDebezium 官方 JDBC Sink 实现topicscounts或counts, types订阅的 Kafka topic多个用逗号分隔transforms.unwrap.typeio.debezium.transforms.ExtractNewRecordState将 Debezium 的before/after结构展开为纯after记录便于 JDBC 写入transforms.unwrap.drop.tombstonesfalse保留 tombstone删除墓碑消息配合delete.enabled支持删除语义insert.modeupsert按主键执行 upsert实现幂等写入delete.enabledtrue允许将 Debezium 删除事件翻译为DELETE语句auto.createtrue/false是否自动建表PG 侧已预建故关闭MySQL 侧开启schema.evolutionbasic开启基础 schema 演进应对字段变更primary.key.moderecord_key主键取自消息的 key即 RisingWave Sink 声明的primary_key idprimary.key.fieldsid指定作为主键的字段名依赖primary.key.moderecord_key上游 RisingWave 建 Sink 时必须显式声明primary_key这也是 create_sink.sql 中primary_key id的用意所在。五、Flink SQL 管道以 Debezium JSON 二次消费管道一由 flink.sql 实现逻辑非常精简——建两张表并做一次 INSERT① Kafka 源表读取 RisingWave 写入的 Debezium 消息CREATE TABLE topic_counts ( id BIGINT, sum BIGINT ) WITH ( connector kafka, topic counts, properties.bootstrap.servers message_queue:29092, properties.group.id test-flink-client-1, scan.startup.mode earliest-offset, format debezium-json, debezium-json.schema-include true );② PostgreSQL 目标表JDBC 连接器主键NOT ENFORCEDCREATE TABLE pg_counts ( id BIGINT, sum BIGINT, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:postgresql://postgres:5432/mydb?usermyuserpassword123456, table-name flinkcounts );③ 执行写入INSERT INTO pg_counts SELECT * FROM topic_counts;两个值得注意的细节format debezium-json且debezium-json.schema-include trueFlink 的debezium-jsonformat 与 RisingWave 的FORMAT DEBEZIUM一一对应schema 包含标志开启后 Flink 可精确解析字段类型scan.startup.mode earliest-offsetFlink 从 topic 最早 offset 开始消费确保测试开始时已有的历史数据也不会遗漏最终可通过行数对比完成端到端校验。六、结果验证一键检查落库数据cargo make setup完成后按 README 可执行以下任一方式验证手动逐表检查# 1. Flink 写入的 PostgreSQL 表 docker exec postgres bash -c psql -U $POSTGRES_USER $POSTGRES_DB -c select * from flinkcounts # 2. JDBC Connector 写入的 PostgreSQL 表 docker exec postgres bash -c psql -U $POSTGRES_USER $POSTGRES_DB -c select * from counts # 3. JDBC Connector 写入的 MySQL 表 docker exec mysql bash -c mysql -u $MYSQL_USER -p$MYSQL_PASSWORD mydb -e select * from counts或直接运行汇总任务cargo make check-db它内部按顺序执行check-pg、check-mysql、check-flink-pg三个子任务。此外还有自动化校验脚本 sink_check.py先等待 60 秒保障 Flink 管道完成消费再对 PostgreSQL 中的counts、flinkcounts、types、flink_types四张表统计行数与预期基线对比并报告失败项适合纳入 CI 作为断言。七、类型兼容性回归全类型极值链路该测试不仅验证基础聚合数据的同步还覆盖了 RisingWave → Kafka → 下游的全类型兼容性。compatibility-rw.sql 在 RisingWave 侧定义了一张覆盖 16 种类型的表kafka_all_typesboolean、smallint、integer、bigint、decimal、real、double precision、varchar、bytea、date、time、timestamp、timestamptz、interval、jsonb 等为其创建types_sink同样FORMAT DEBEZIUM ENCODE JSONtopictypes随后插入三行极值数据各数值类型的上下限如-9223372036854775807、9223372036854775807超长 varchar 与 bytea 大字段日期时间极值0001-01-01、9999-12-31 23:59:59interval 9990 year等边界区间。这些数据经 Kafka 到达 PG 的types表JDBC Connector 写入与flink_types表Flink 写入。由于 PG 预置表中c_decimal、c_interval以text类型承载正好规避了 Debezium JSON 数值精度与 interval 语法的跨引擎差异体现了一种务实的兼容性验证手法。验证命令同样收敛在 Makefile 中cargo make check-com-pg # select * from types cargo make check-flink-com-pg # select * from flink_types八、原理小结与可复用要点从源码结构与任务编排可以推断该测试的价值在于验证了以下事实链格式契约RisingWave 的FORMAT DEBEZIUM ENCODE JSON产出的是标准 Debezium 消息能被 Flinkdebezium-jsonformat 与 Debezium Connect 原生解析说明 RisingWave 的 Kafka Sink 具备与 Debezium 生态直接互操作的能力主键语义贯穿全链路RisingWave 侧primary_key→ Kafka 消息 key → Debezium Connectprimary.key.moderecord_key→ 目标表主键一环缺失都会导致 upsert 失效双消费端并行验证同一条 Kafka topic 可被流处理框架Flink与连接器框架Debezium JDBC Sink同时消费互不干扰验证了消息协议的通用性类型与极值回归通过全类型表和边界值验证 Debezium JSON 序列化在极端数据下的正确性。如需在自有环境中复现只需确保cargo make与 Docker 可用然后在本仓库根目录执行cargo make setup任务定义见 Makefile.toml完成后用cargo make check-db与cargo make check-com-pg/cargo make check-flink-com-pg验证两条管道与类型兼容性链路即可。赞分享数据库流处理后端数据工程【免费下载链接】risingwaveEvent streaming platform for agentic AI. Continuously ingest, transform, and serve event streams in real time, at scale.项目地址https://gitcode.com/gh_mirrors/ri/risingwave点击查看免费下载相关推荐Debezium JDBC Sink Connector 的 DB2iIBM i / AS/400Dialect测试策略与 SQL 生成实现解析Debezium JDBC Sink Connector 的 DB2iIBM i / AS/400Dialect测试策略与 SQL 生成实现解析 本文以后端变更数据捕获数据集成流处理SeaTunnel Kafka 连接器 COMPATIBLE_KAFKA_CONNECT_JSON 格式实战消费 Kafka Connect JDBC 与 Debezium 数据SeaTunnel Kafka 连接器 COMPATIBLE_KAFKA_CONNECT_JSON 格式实战消费 Kafka Connect JDBC 与 D数据集成ETL大数据批处理流处理变更数据捕获Flink CDC 实战从零搭建 MySQL 到 Kafka 的流式 ELT 管道Flink 2.2 快速上手Flink CDC 实战从零搭建 MySQL 到 Kafka 的流式 ELT 管道Flink 2.2 快速上手 本文基于 Flink CDC 官方文档中的后端数据集成大数据流处理变更数据捕获数据同步上一篇突破传统Java APNS如何重塑你的iOS推送开发体验下一篇从源码解析Material-ish ProgressProgressWheel类的设计与实现细节创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于恒美微站

恒美微站专注于为个体商户、工作室提供极简自助建站服务,让每个人都能轻松拥有专业网站。

快速链接

  • 关于我们
  • 建站服务
  • 主题模板
  • 案例展示
  • 资讯中心

服务项目

  • 可视化建站
  • 拖拽编辑
  • 主题定制
  • SEO 优化
  • 网站托管

联系方式

  • 📍 地址:北京市朝阳区建国路 88 号
  • 📞 电话:400-888-8888
  • ✉️ 邮箱:info@hmyw.cn
  • 🕐 时间:周一至周日 9:00-18:00

© 2024 恒美微站 hmyw.cn 版权所有 | 京 ICP 备 12345678 号