ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Paimon 数据出仓源码导读(四):从零跑通 Paimon 到 MySQL:环境、建表与第一条数据

Apache Paimon 数据出仓源码导读(四):从零跑通 Paimon 到 MySQL:环境、建表与第一条数据 上一篇我们从源码角度总览了 Paimon 出仓到 MySQL 时的 UPSERT、DELETE、批量写入和故障重放。不过如果还没有亲手跑过这条链路一上来就看RowKind、JdbcOutputFormat、executeBatch()和 Checkpoint很容易把几个层次混在一起。所以从这一篇开始我们暂时放慢速度把“Paimon 出仓到 MySQL”拆成几篇独立文章。这一篇只完成一件事从零搭好一条最小链路亲眼看到 Paimon 中的一条订单在 MySQL 中新增、更新和删除。先不追每一个源码调用也先不讨论失败重放。等链路真正跑通以后下一篇再把“一条记录怎样进入 JDBC Sink、怎样在 Buffer 中攒成 Batch”逐层拆开。一、这次到底要搭一条什么链路最终运行的是一条持续存在的 Flink SQL 作业Paimon 主键表 orders ↓ Paimon Source 持续读取当前全量和后续变化 ↓ Flink SQLINSERT INTO mysql_orders SELECT ... ↓ JDBC Sink ↓ MySQL 物理表 orders_rt这里有三张“表”但它们并不都是同一种东西名字存在哪里它是什么ordersPaimon真正保存订单数据的湖表mysql_ordersFlink Catalog指向 MySQL 的 JDBC 表定义本身不保存另一份数据orders_rtMySQL真正接收出仓结果的物理表很多初学者会误以为执行CREATE TABLE mysql_orders后Flink 会自动去 MySQL 建表。实际上它只是向 Flink 注册了一份连接和字段映射。MySQL 中的orders_rt仍然需要提前创建。二、运行前需要准备哪些组件本文环境基于Apache Paimon 1.4.2 Apache Flink 1.20.1 Flink JDBC Connector 3.3.0-1.20 MySQL 8.x除了 Flink 和 Paimon还需要两个容易漏掉的依赖依赖解决什么问题缺失时常见现象Flink JDBC Connector让 Flink 认识connector jdbc找不到jdbcFactoryMySQL Connector/J让 JDBC 能真正连接 MySQL找不到 MySQL Driver 或无法创建连接Flink 的二进制发行包通常不自带 JDBC Connector 和数据库 Driver。对于普通 Standalone 集群最直观的做法是让 JobManager 和所有 TaskManager 的 Flinklib目录都包含版本匹配的 JAR然后重启对应进程。为什么不能只放在 SQL Client 所在机器因为 SQL Client 主要负责解析和提交作业真正执行 JDBC Sink 的是 TaskManager。TaskManager 看不到 Driver作业仍然会失败。先做三个检查第一确认 Paimon 和 Flink 大版本匹配。不要把面向 Flink 1.18 的 Bundle 直接放进 Flink 1.20。第二确认 JDBC Connector 的后缀与 Flink 版本匹配。本文使用的是3.3.0-1.20。第三确认 MySQL 网络可达。TaskManager 所在机器必须能访问 MySQL Host 和端口不是只有你的电脑能连上就够了。三、先在 MySQL 建真正的目标表先创建实验数据库CREATEDATABASEIFNOTEXISTSservingDEFAULTCHARACTERSETutf8mb4;再创建目标表CREATETABLEserving.orders_rt(idBIGINTNOTNULL,amountDECIMAL(10,2),statusVARCHAR(32),PRIMARYKEY(id))ENGINEInnoDB;这条PRIMARY KEY (id)不是可有可无的装饰。后面同一个订单再次到来时MySQL 正是依靠这个真实主键判断id 不存在 - 插入新行 id 已存在 - 更新原来的行创建完成后不要只凭印象判断主键已经存在直接检查SHOWCREATETABLEserving.orders_rt;结果中应该能看到类似PRIMARYKEY(id)写入账号需要哪些权限JDBC Sink 至少会执行 INSERT、UPDATE 语义和 DELETE因此写入账号要具备目标表对应权限。生产环境中建议由 DBA 创建最小权限账号不要为了省事给全库管理员权限。本文 DDL 中继续使用占位账号username flink_writer password ******不要把真实密码直接提交到 Git、文章截图或公开的 Flink SQL 文件中。四、在 Flink SQL Client 中准备实验环境进入 Flink SQL Client 后先使用流模式SETexecution.runtime-modestreaming;为了让第一次实验更容易观察可以暂时把默认并行度设为 1并开启 10 秒 CheckpointSETparallelism.default1;SETexecution.checkpointing.interval10s;并行度设为 1 只是为了让初学实验中的日志、连接和数据顺序更容易看懂不是生产推荐值。生产环境的并行度要结合 MySQL 连接数、写入吞吐、热点 Key、锁等待和 Checkpoint 时长重新评估。五、创建 Paimon 订单主键表在 Flink SQL 中创建 Paimon 表CREATETABLEorders(idBIGINT,amountDECIMAL(10,2),statusSTRING,PRIMARYKEY(id)NOTENFORCED)WITH(connectorpaimon,pathhdfs:///warehouse/orders,changelog-producerinput);逐项解释定义含义id订单唯一标识amount订单金额保留两位小数status订单状态PRIMARY KEY (id)相同id表示同一份订单状态NOT ENFORCEDFlink 不负责运行时检查唯一性pathPaimon 表在文件系统中的存储位置changelog-producer input保存上游输入变化供下游增量消费这里先使用input是为了后面能更直观地观察变化怎样传到 JDBC Sink。不同 Changelog Producer 的差别已经在上一篇单独讲过。先写两条初始数据执行INSERTINTOordersVALUES(1001,CAST(80.00ASDECIMAL(10,2)),CREATED),(1002,CAST(50.00ASDECIMAL(10,2)),CREATED);作业完成后可以在用于查询快照的 SQL Client 会话中切到 Batch 模式再查询 Paimon 当前结果SETexecution.runtime-modebatch;SELECTid,amount,statusFROMordersORDERBYid;预期当前结果是idamountstatus100180.00CREATED100250.00CREATED这里看到的是 Paimon 表的当前状态不是历史变化流水。六、在 Flink 中注册 MySQL JDBC Sink现在创建 Flink JDBC 表CREATETABLEmysql_orders(idBIGINT,amountDECIMAL(10,2),statusSTRING,PRIMARYKEY(id)NOTENFORCED)WITH(connectorjdbc,urljdbc:mysql://mysql-host:3306/serving,table-nameorders_rt,usernameflink_writer,password******,sink.buffer-flush.max-rows100,sink.buffer-flush.interval1s,sink.max-retries3);这一段最容易产生两个误解。误解一mysql_orders是又创建了一张 MySQL 表不是。mysql_orders是 Flink 中的逻辑表名table-name orders_rt才是 MySQL 里的真实物理表。以后在 Flink SQL 中写INSERTINTOmysql_orders...Connector 才会根据 URL 和table-name找到serving.orders_rt误解二Flink DDL 的主键会替 MySQL 创建约束也不会。Flink DDL 中的PRIMARYKEY(id)NOTENFORCED告诉 Planner 和 JDBC Connector这是一张按id更新和删除的 Upsert Sink。MySQL DDL 中的PRIMARYKEY(id)才是数据库真正执行的唯一约束。两边都要有而且 Key 语义必须一致。七、第一次启动 Paimon 到 MySQL 的出仓作业在用于运行同步作业的 SQL Client 会话中确认切回 Streaming 模式再执行SETexecution.runtime-modestreaming;INSERTINTOmysql_ordersSELECTid,amount,statusFROMorders/* OPTIONS( scan.mode latest-full, consumer-id mysql-orders-v1 ) */;这不是执行完马上退出的普通查询而是一条持续运行的流作业。latest-full可以先按下面这句话理解启动时读取最新 Snapshot 的当前全量之后继续读取新产生的变化。所以第一次启动时前面已经存在的1001和1002也会进入 MySQL不需要等它们再次发生变化。如果 SQL Client 以 Attached 模式运行这个终端可能会一直被作业占用。后续写入 Paimon 和执行 UPDATE、DELETE 时可以再开一个 SQL Client 会话。八、先验证启动时全量是否写进 MySQL在 MySQL 中执行SELECTid,amount,statusFROMserving.orders_rtORDERBYid;等待 JDBC Sink Flush 后预期看到idamountstatus100180.00CREATED100250.00CREATED如果暂时查不到不要第一时间判断数据丢了。本文配置了sink.buffer-flush.interval1s低流量时记录可能先在 JDBC Sink Buffer 中等待下一次定时 Flush。除此之外还要看作业是否处于 RUNNING、Checkpoint 是否正常、TaskManager 日志有没有连接错误。九、再写一条新订单验证增量新增在另一个 Flink SQL Client 会话中执行INSERTINTOordersVALUES(1003,CAST(120.00ASDECIMAL(10,2)),CREATED);Paimon 会产生新的 Snapshot。持续运行的 Source 发现它以后把新订单交给 JDBC Sink。再查 MySQLSELECTid,amount,statusFROMserving.orders_rtWHEREid1003;预期结果idamountstatus1003120.00CREATED到这里我们已经验证了两种读取作业启动前存在的 1001、1002 - 启动全量 作业启动后新增的 1003 - 持续增量十、同一个主键再写一次验证更新现在让订单1001从 80 元变成 100 元状态变成PAID。对 Paimon Deduplicate 主键表可以再次写入完整的新值INSERTINTOordersVALUES(1001,CAST(100.00ASDECIMAL(10,2)),PAID);虽然 SQL 写的是INSERT INTO但id1001已经存在。对当前表状态来说这是同一主键的新版本不应该再多出第二个1001。查询 PaimonSELECTid,amount,statusFROMordersWHEREid1001;预期只有一行idamountstatus1001100.00PAID再查询 MySQLSELECTid,amount,statusFROMserving.orders_rtWHEREid1001;MySQL 也应该仍然只有一行并且金额和状态都已经更新。为什么 MySQL 没有插出两行JDBC Sink 对主键表使用 MySQL UPSERT核心语句类似INSERTINTOorders_rt(id,amount,status)VALUES(?,?,?)ONDUPLICATEKEYUPDATEidVALUES(id),amountVALUES(amount),statusVALUES(status);第一次id1001不存在走 INSERT。第二次id1001已存在MySQL 主键发生重复走 UPDATE 分支。所以“更新订单”不代表 Connector 一定发送普通的UPDATEorders_rtSET...WHEREid...;对 MySQL JDBC Upsert Sink新增和更新通常共用INSERT ... ON DUPLICATE KEY UPDATE。能不能直接执行 UPDATE Paimon 表Paimon 1.4.2 在 Flink 1.17 及以上支持对主键表执行 UPDATE但 UPDATE 是 Batch 模式操作而且不能修改主键。可以在另一个用于批处理的 SQL Client 会话中执行SETexecution.runtime-modebatch;UPDATEordersSETamountCAST(110.00ASDECIMAL(10,2)),statusPAIDWHEREid1001;这次 Batch DML 提交新的 Paimon Snapshot 后原来持续运行的出仓作业仍会继续发现并同步变化。十一、删除订单验证 DELETEPaimon 的DELETE FROM同样在 Batch 模式执行并且只支持满足条件的主键表与 Merge Engine。在批处理 SQL Client 会话中执行SETexecution.runtime-modebatch;DELETEFROMordersWHEREid1001;删除完成后先查 Paimon 当前状态SELECTid,amount,statusFROMordersORDERBYid;1001应该已经不存在。持续运行的 Source 读到 Delete Changelog 后JDBC Sink 最终会按主键执行类似DELETEFROMorders_rtWHEREid?;再查 MySQLSELECTid,amount,statusFROMserving.orders_rtWHEREid1001;预期返回 0 行。十二、把订单1001的完整过程串起来现在不看源码只看业务状态时刻Paimon 操作Paimon 当前结果MySQL 当前结果T1首次写入1001, 80, CREATED1001, 80, CREATED出仓后相同T2再写1001, 100, PAID1001, 100, PAIDUPSERT 后相同T3DELETE1001不存在DELETE 后不存在Paimon 和 MySQL 之间传递的不是“每隔一段时间复制整张表”而是一条持续运行的变化链路。从最终状态看新增MySQL 出现一行 更新同一主键的值被覆盖行数不增加 删除MySQL 中对应主键消失至于 T1 到 T2 之间究竟发出了U还是-D/I它们有没有在同一个 Buffer 中合并什么时候调用executeBatch()留到下一篇继续拆。十三、两个主键为什么缺一不可把两边的定义放在一起-- Flink JDBC Sink DDLPRIMARYKEY(id)NOTENFORCED-- MySQL 物理表 DDLPRIMARYKEY(id)它们分别解决不同问题主键使用者作用Flink DDL 主键Planner、JDBC Connector识别 Upsert Key允许处理 UPDATE 和 DELETEMySQL 物理主键MySQL阻止重复 Key触发 UPSERT 的 UPDATE 分支只有 Flink 主键没有 MySQL 主键Connector 仍可能生成 UPSERT SQL但 MySQL 找不到重复键冲突。同一个id再来一次时目标表可能出现重复行幂等恢复也失去基础。只有 MySQL 主键没有 Flink 主键Flink 会把 Sink 当成 Append 模式。上游查询一旦包含 UPDATE 或 DELETE规划阶段通常就会拒绝或者无法得到期望的更新语义。两边都有主键但字段不一致如果 Paimon 用(tenant_id, order_id)标识订单MySQL 却只用order_id不同租户的相同订单号会互相覆盖。因此主键检查不能只看“都有 PRIMARY KEY”还要检查字段数量、顺序、类型和业务语义。十四、latest-full和consumer-id分别在做什么本文 Source 使用scan.modelatest-full,consumer-idmysql-orders-v1latest-full它解决第一次启动从哪里读先读取最新 Snapshot 的完整当前状态 再持续读取后续新变化这适合第一次为一张空的 MySQL 服务表建立镜像。consumer-id它给这一路长期消费一个稳定身份Paimon 可以据此管理 Consumer 相关进度和 Snapshot 保留。但要注意consumer-id不是 Flink Checkpoint也不能单独保证作业故障后精确恢复到某一条记录。Flink Source Split、读取位置和算子状态仍然需要 Checkpoint 或 Savepoint 保存。十五、为什么刚写入后 MySQL 可能还查不到本文配置sink.buffer-flush.max-rows100,sink.buffer-flush.interval1s记录到达 JDBC Sink 后不一定立即访问 MySQL。它可能先进入当前 Sink 子任务的内存 Buffer直到下面任一条件发生收到的记录达到 100 条 等待时间达到 1 秒 Flink 开始做 Checkpoint 作业正常关闭并清理尾批所以低流量实验中看到约 1 秒的可见延迟通常是正常现象。如果长时间仍没有数据再检查Flink 作业是不是 RUNNINGSource 有没有读到新 SnapshotSink 有没有持续报 JDBC 异常MySQL 是否存在锁等待或连接耗尽Checkpoint 是否频繁失败查询的是不是正确的数据库和表。十六、第一次跑最常见的八类错误1. 找不到 JDBC Factory典型信息包含Could not find any factory for identifier jdbc优先检查 Flink JDBC Connector JAR 是否存在、版本是否匹配、集群进程是否已经重启。2. 找不到 MySQL Driver典型信息包含No suitable driver ClassNotFoundException: com.mysql.cj.jdbc.Driver优先检查 MySQL Connector/J 是否在真正执行作业的 TaskManager Classpath 中。3. MySQL 拒绝连接常见原因包括 Host、端口、账号、密码、授权来源 Host、防火墙和 TLS 配置不正确。不要只在本机测试 MySQL Client要从 TaskManager 所在网络环境验证可达性。4. Flink 中注册了主键MySQL 却出现重复数据执行SHOWCREATETABLEserving.orders_rt;确认 MySQL 物理表真的存在 PRIMARY KEY 或语义完全一致的 UNIQUE KEY。5. 新增能写更新或删除规划失败检查 Flink JDBC Sink DDL 是否声明主键以及上游查询经过投影、Join、聚合后是否仍保留可用的 Upsert Key。6. DELETE 执行后 MySQL 仍有数据先确认 Paimon 的 DELETE DML 是否成功提交新 Snapshot再确认 Changelog Producer 是否产生删除变化最后检查两边 Key 是否一致。7. 金额或字符串写入失败检查 Paimon、Flink JDBC DDL 和 MySQL 三边的数据类型。例如Paimon DECIMAL(10, 2) Flink Sink DECIMAL(10, 2) MySQL DECIMAL(10, 2)字段名字相同不代表类型一定兼容。精度、长度、NULL 约束和字符集都可能导致失败。8. 作业正常但目标端看起来延迟很大先看sink.buffer-flush.interval再看 Sink 是否背压、MySQL 是否慢、Checkpoint 是否长时间执行。十七、跑通以后做一次最小验收不要只看到一条 INSERT 成功就宣布链路完成。至少验证下面六项验收项期望结果启动全量作业启动前的 Paimon 当前数据进入 MySQL增量新增新 Key 出现在 MySQL同 Key 更新MySQL 仍只有一行字段变为新值删除MySQL 对应 Key 消失NULL 和边界值类型、长度、精度符合预期Checkpoint能持续成功不只是作业显示 RUNNING再补一项非常实用的核对-- Paimon 侧SELECTCOUNT(*)FROMorders;-- MySQL 侧SELECTCOUNT(*)FROMserving.orders_rt;行数相同只是第一步不足以证明内容完全一致。正式上线还要对主键集合、关键字段、删除和抽样明细做对账。十八、这一篇先记住五句话1. mysql_orders 是 Flink 逻辑表orders_rt 才是 MySQL 物理表 2. Flink DDL 主键决定 Upsert 语义MySQL 主键真正执行唯一约束 3. latest-full 先读当前全量再持续读后续变化 4. 新增和更新通常通过 MySQL UPSERT 落地删除按主键执行 DELETE 5. 写入先进入 JDBC Sink Buffer低流量时不一定立刻在 MySQL 可见到这里我们已经从零跑通了一条最小的 Paimon 到 MySQL 出仓链路。下一篇不再停留在 DDL 层而是拿四条真实变化逐步跟进源码Paimon Source 怎样一条条发 RowData GenericJdbcSinkFunction.invoke() 为什么每次只收一条 TableBufferReducedStatementExecutor 怎样按主键覆盖 为什么收到 100 条输入最终不一定执行 100 条 MySQL DML addBatch() 和 executeBatch() 到底分别做了什么把这段过程看清楚以后“出仓到 MySQL 到底是一条一条还是一批一批”就不再只是一句结论而是一条能从源码、日志和 MySQL 现象互相验证的完整链路。本篇关键配置与资料位置docs/content/flink/sql-write.mdPaimon 的 INSERT、UPDATE、DELETE DML 说明JdbcDynamicTableSink.javaFlink JDBC Sink 的 ChangelogMode 和主键校验MySqlDialect.java生成 MySQLINSERT ... ON DUPLICATE KEY UPDATEJdbcConnectorOptions.java定义 Buffer Flush、时间间隔和重试配置Flink 1.20 JDBC Connector 文档依赖、Upsert 模式、主键要求和 Connector Options本文基于 Apache Paimon 1.4.2、Apache Flink 1.20.1 和 Flink JDBC Connector 3.3.0-1.20。不同版本的依赖坐标、默认值和类名可能变化实际部署时请以对应版本文档和源码为准。
返回列表