尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

ClickHouse 25.12新特性与Flink CDC实时同步实践

ClickHouse 25.12新特性与Flink CDC实时同步实践 1. 25.12 更新了什么三个让我觉得“这次没白等”的变化先说说这次发布说明里最抓眼球的三件事。第一个是Pipeline执行模型默认开启第二个是向量化窗口函数落地第三个是零拷贝复制Zero Copy Replication终于生产可用顺带把NFS依赖彻底拿掉了。这三点放在一起看基本能判断出ClickHouse现在的走向单查询朝着更细粒度的并行调度去优化分析函数朝着能处理复杂事件序列的方向去补集群运维则朝着去掉外部依赖、降低托管成本的方向去收敛。Pipeline执行模型这事我最早是在21.x版本开始关注当时还是实验特性要用set allow_experimental_pipeline_model1才能打开。这个模型的本质是把一次查询拆成很多个小任务塞进一个类似流式处理引擎的调度框架里而不是像旧模型那样按照“读取数据—过滤—聚合”这种粗粒度阶段一层层往下压。我自己在测试环境跑过好几轮对比比如对一张20亿行的宽表跑GROUP BY city, count()开启Pipeline之后查询尾延迟的抖动明显变小CPU的利用率曲线也更平滑。25.12把它默认打开意味着普通用户不需要再关心这些底层开关升级之后如果没特殊配置跑的就是新模型。向量化窗口函数是我个人最期待的一个。过去在ClickHouse里写row_number()、lag()这类窗口函数虽然也能跑但实现路径偏慢尤其在大表上排序开窗经常被吐槽“能用但别太指望性能”。25.12把窗口函数执行重写成了向量化版本我测试过一个实际场景一张订单明细表要按用户ID分组、按下单时间排序取每个用户最近三笔订单。旧版本跑一次大概要40秒25.12上同一套SQL跑到了12秒左右。这个提升幅度对于写实时报表、做留存分析、算用户行为序列的团队来说属于“升级后SQL不用改性能白捡”的福利。零拷贝复制这块如果你维护过大规模ClickHouse集群应该知道以往ReplicatedMergeTree表的副本同步依赖ZooKeeper做元数据协调而且数据文件依然通过本地磁盘复制来保持一致。25.12里零拷贝复制被标记为生产可用加上NFS依赖移除进一步印证了ClickHouse正在把“跨副本复制”从“文件级复制”推向“元数据级复制”。生产环境下这意味着什么副本数量增多时网络开销和磁盘占用能明显下降日常运维也不必再为一个NFS挂载点折腾权限、容灾和IO延迟。不过我要提醒一句零拷贝复制在设计上假设“不同副本上的同名Part是等价的”所以任何绕过MergeTree引擎直接往数据目录里丢文件、或者用remote()表函数跨副本散写的方式都容易破坏这个假设升级前最好审查一下有没有这类“野路子”操作。这次发布还有一个容易被忽略的细节查询参数化能力进一步增强。我理解ClickHouse团队是想让网关层、BI工具层能更好地复用服务端预编译的查询计划。对于用Metabase、Superset或者自研查询平台接ClickHouse的团队这个改进能减少每次请求的解析开销间接提升高并发查询的吞吐。2. 升级前必须检查的清单从兼容性到配置项变化升级ClickHouse这种列式数据库最怕的不是版本号变了而是“行为变了但没人通知你”。尤其是大版本跨越比如从22.x或者23.x直接跳到25.12中间积累了大量的默认行为调整。这里我按自己生产环境升级的检查顺序整理一份通用清单。第一项确认ZooKeeper/Keeper版本兼容性。25.12在使用ReplicatedMergeTree时仍然依赖Keeper或ZooKeeper做协调虽然零拷贝复制降低了数据复制压力但元数据协调链路没有消失。如果集群里的ClickHouse Server还是3年以前的老版本配套的ZooKeeper可能也停留在3.4、3.5这种古老版本。建议升级前把ZK至少升到3.7以上或者直接用ClickHouse Keeper替代。不然后续建表、分区卸载这些操作会因为协议兼容问题出现莫名其妙的超时。第二项检查allow_experimental_*这类开关是否还被识别。25.12把不少原本需要手动开启的实验特性转正了例如Pipeline模型。如果你在旧的config.xml或users.xml里写死了allow_experimental_pipeline_model1升级后这个参数可能不再被识别。ClickHouse对“不存在的配置项”通常只给一个warning不会直接拒绝启动但它会造成“我以为开了其实没开”的错觉。我在测试升级时就遇到过类似情况旧配置里一个实验开关在24.x版本已被移除服务启动后日志里只是轻描淡写地提示“Unknown setting”结果一些查询路径又回到了老实现。升级后建议打开system.settings表核对一下计划使用的参数是否还在生效。第三项MergeTree系列的min_bytes_for_full_part、max_bytes_to_merge_at_min_space_in_pool等后台参数默认值可能变了。25.12对后台merge策略做了一些调整目的是减少小Part数量、降低写入放大。我在测试实例上观察到的现象是升级后同样一批写入系统自动触发的merge更激进Part数量比老版本少大概20%到30%。这对查询性能是好事但对磁盘空间的使用会更敏感。如果你用了TTL做数据过期清理而且要依赖Part级别的删除逻辑那么merge策略变化会导致TTL的实际执行粒度发生偏移。经验做法是升级后在system.parts里连续观察一周确认Part数量和大小分布没有异常波动。第四项查询语法兼容性回归。我整理过一份常见的不兼容点表格这里直接贴出来旧写法≤23.x25.12的表现建议改法ARRAY JOIN与LEFT ARRAY JOIN混用不加括号部分语义被重新解析返回行数可能不同显式加括号区分嵌套WITH FILL搭配STEP但类型是字符串报类型错误先用toDateTime等函数转换JOIN时USING字段与SELECT字段混淆审查更严格可能报“Ambiguous column”显式写t1.col、t2.col自定义函数SQL UDF使用CREATE FUNCTION行为保留但参数名大小写敏感度调整统一小写参数名这一步没有捷径只能把线上跑的核心查询脚本、报表SQL在测试环境跑一遍回归。别信“版本兼容性列表”业务SQL千奇百怪只有实际执行过才知道坑在哪。第五项磁盘与内存配置的侧重点变化。25.12的Pipeline模型在调度上更积极对内存池的分配粒度也做了调整。如果你之前的max_memory_usage设置是基于老模型压测出来的阈值升级后可能需要重新压测。我遇到过的情况是同样一个复杂分析查询老版本内存峰值60GB25.12跑到55GB但CPU使用率更均匀。这里没有“一定变好”的说法只有“重新测一下才知道”。3. 把Flink与ClickHouse 25.12结合起来MySQL实时同步到ClickHouse的完整方案说完了版本本身我打算花大篇幅讲讲如何用Flink把MySQL的数据实时同步到ClickHouse。这不是空穴来风——在实际大促、报表、用户画像场景里MySQL负责在线交易ClickHouse负责分析查询中间隔着一道“数据管道”而Flink是当前最适合承担这道管道的计算引擎。在选择这个方案之前我对比过几种常见工具比如Canal监听Binlog然后写到Kafka再进ClickHouse或者用DataX做离线批量同步再或者直接用ClickHouse的mysql表引擎做实时查询。Canal那套链路适合“已经有Kafka基础设施、且下游有多个消费方”的团队DataX只适合T-1离线同步mysql表引擎则是把查询下推到MySQL分析复杂一点MySQL就扛不住。Flink CDC的优势在于既能读MySQL全量快照也能无缝切换增量Binlog还自带Checkpoint、状态管理和流式SQL和ClickHouse对接的成熟度这几年提升得很明显。下面我按“全量同步、增量同步、DDL变更、目标表设计、性能调优、常见坑”六个部分展开。3.1 环境准备Flink与ClickHouse的版本搭配先说版本组合。我的测试环境如下ClickHouse 25.12.1社区版部署方式为单机测试 三节点集群Flink 1.18.1Flink CDC 3.1.1这个版本对MySQL和StarRocks支持很好和ClickHouse搭配也没问题JDBC驱动com.clickhouse:clickhouse-jdbc:0.6.3MySQL 8.0.32开启Binlogbinlog_row_imageFULL如果你用Flink 1.17或1.19问题也不大Flink CDC的连接器对Flink版本有一定兼容要求建议看下对应版本的官方文档。ClickHouse JDBC驱动要选0.5.0以上不然对DateTime64和Decimal类型的映射会出问题。在系统层面需要给Flink TaskManager适当调大堆内存。同步任务看起来是“搬数据”但Flink在解析Binlog、序列化Row、维护状态时很吃内存。我在一个节点上分配了8GB堆外内存最终稳定在5GB左右。3.2 全量同步阶段Flink CDC的Snapshot机制怎么用Flink CDC在做全量同步时并不像DataX那样直连查一遍。它的逻辑是先获取一个全局读锁如果表不是空的然后读取Binlog位点再启动“快照读取线程”扫描表数据扫完后释放读锁再接着从之前记录的Binlog位点消费增量。这套机制的好处是全量数据和增量数据之间不会漏数据坏处是如果MySQL表特别大快照阶段持锁时间长会对线上写入产生阻塞。我自己的经验是控制以下几点。分批拉取Flink CDC的SnapshotReader内部其实会按主键范围分段拉取。只要表有主键它默认会切成多个chunk并行读取。如果表没有主键Flink CDC会退化为“全表单线程扫描”速度会慢一个量级。所以目标同步的表必须有主键没有主键就先加一列自增ID或者用联合主键。调整chunk大小默认的chunk大小是8096行。对MySQL来说如果单行字段特别宽比如有TEXT、JSON可以调小到4096以免一次拉取的数据过大导致网络包堆积。如果单行很小且表数据在千万行以下可以调到16384减少拆分次数。并行度设置全量阶段建议把source并行度设为4到8我测试时设为6。注意Flink CDC是全量快照阶段并行增量阶段由于要保证Binlog顺序并行度有限制所以不用在这个阶段把并行度调到32之类的极端值。Checkpoint全量阶段也要开启Checkpoint。Flink CDC会把已经完成的chunk信息存到状态里如果任务中途挂了可以从Checkpoint恢复已经同步的chunk不需要重新扫这在数据量大的场景里能省不少时间。我遇到过全量扫500GB表跑到57%时网络抖动导致TaskManager失联好在Checkpoint间隔设置的是60秒恢复后从最近一次Checkpoint继续没有从头再来。全量同步完成后Flink CDC会自动“感知”到当前Binlog位点随后进入增量模式不需要人工干预。这个切换在UI上能看到source的指标从“Snapshot Phase”变为“Binlog Phase”。3.3 增量同步阶段Binlog解析和Exactly-Once语义进入增量阶段后Flink CDC读取MySQL的Binlog解析成一个个变更事件这里有一个对ClickHouse场景很关键的细节ClickHouse不是OLTP数据库它不擅长高频小型写入。MySQL里一张表可能有每秒几百上千次的UPDATE和DELETE如果这些变更都一条条直接写入ClickHouse效果会非常差。原因有两个一是ClickHouse单次插入的优化对象是“一批数据”几百条一批对MergeTree来说也算小批次会产生大量小Part二是ClickHouse的UPDATE和DELETE是异步Mutation频繁执行会在后台制造大量重写操作严重拖慢merge线程。所以我的做法是引入窗口聚合缓冲把Mini-Batch写入变成真正的批量写入。操作方式不复杂。用Flink Table API或DataStream在sink前面加一个自定义的buffer算子攒满2000条或超过3秒触发一次写入。写入时按ClickHouse的INSERT INTO ... VALUES或INSERT INTO ... SELECT格式批量拼接。实测下来一批2000条的写入和一条条写入相比ClickHouse端Part数量减少了一个数量级写入吞吐提升了大概6到10倍。这里还有一道坎是Exactly-Once。Flink CDC从MySQL读到数据如果要实现不丢不重需要sink支持事务或幂等。ClickHouse的JDBC连接器本身对事务的支持比较弱尤其是MergeTree表官方建议的方式是“用异步写入以幂等为目标”。我的取舍是这样的核心交易类指标允许少量重复靠ClickHouse端的查询去重兜底比如用argMax取最新值。日志类、行为类数据允许重复但不允许丢失靠Checkpoint保证“至少一次”。真正要求精密的场景给目标表加ReplacingMergeTree引擎以业务主键做去重。说实话ClickHouse和Flink之间要做到完美Exactly-Once目前仍然要结合业务场景做取舍。如果读者正在选型我的建议是不要盲目追求“端到端唯一一次”不如把重复控制通过ClickHouse自身的去重引擎解决这样实现成本低线上稳定性反而更高。3.4 DDL变更同步新增列要小心什么MySQL表结构变了比如新增一列、修改字段长度Flink CDC能否自动同步到ClickHouse答案是能但受限很多。Flink CDC 3.x支持通过include.schema.changestrue配置来传递DDL事件但下游的JDBC sink不一定能自动执行DDL。对ClickHouse我的处理方式是写一个自定义的Sink算子监听SchemaChangeEvent解析出SQL语句后经过白名单校验只允许ALTER TABLE ADD COLUMN拦截DROP COLUMN和MODIFY COLUMN再通过JDBC执行。这件事的复杂度在于MySQL和ClickHouse的类型体系差异。比如MySQL的DATETIME(3)在ClickHouse里可能需要映射为DateTime64(3)MySQL的varchar(255)根据Collation不同字符集可能影响ClickHouse侧的String长度限制MySQL的DECIMAL(10,2)ClickHouse里有Decimal(10,2)但要注意ClickHouse的Decimal整数位上限是76位精度常用范围内没差别。我的做法是维护一张字段映射表像这样MySQL字段类型ClickHouse字段类型注意事项TINYINTInt8注意无符号版本需手动转UInt8INTInt32无符号转UInt32BIGINTInt64无符号转UInt64VARCHAR(n)String不需要指定长度DATETIMEDateTime时区问题需要统一DATETIME(3)DateTime64(3)精度对齐TIMESTAMPDateTime受MySQL时区影响需用connectionTimeZone控制DECIMAL(p,s)Decimal(p,s)注意p不能超过76TEXT/LONGTEXTString无长度限制JSONStringClickHouse侧可用JSONExtract做后续解析在DDL同步的稳定性上我强烈建议先在小表上验证再推广到大表对于大表的DDL变更宁可停机手工处理也不要依赖自动同步。因为ClickHouse的ALTER TABLE操作在Part数量极多时会排队如果Flink端一直重试执行DDL可能拖垮整个同步链路。3.5 目标表设计与写入模式MergeTree、ReplacingMergeTree还是CollapsingMergeTreeClickHouse端目标表怎么建是很多人容易忽略的环节。同步过去的表是复制MySQL的表结构还是重新设计成分析友好的模型我在实践中遵循下面这套原则。流水表日志、订单、行为直接对应MySQL表结构用MergeTree。排序键选用event_time或者其他时间字段加上经常过滤的维度字段。如果MySQL表有自增ID排序键可以设计为(id, event_time)这样按ID查明细很快。维表用户、商品、配置用ReplacingMergeTree。在以MySQL为主的生产系统里这些表随时可能被更新。ReplacingMergeTree基于排序键去重相同排序键的记录只会保留最新一条。配合Flink端按主键去重的upsert逻辑能保证ClickHouse查询时拿到最新状态。状态快照表、余额表、库存表用CollapsingMergeTree。这类表的特点是“同一主键存在多行”每行带一个sign字段1表示新增/更新-1表示取消旧值。Flink端在写入时将上游的一条变更记录拆成“先插入-1取消旧值再插入1新增新值”两条ClickHouse在merge时自动把sign1和sign-1的两行抵消。这样查总数时用sum(sign)聚合拿到的就是当前真实值。此外表的分区策略要结合数据量和查询模式。时间字段按天分区是默认选择但如果你经常按周或按小时查询可以调整为按周或按小时分区。分区过多比如几百上千个会影响SELECT的元数据扫描所以我习惯将控制分区的代码写在同步任务里每次写入时显式指定分区键而不是让ClickHouse按默认规则生成。写入模式上ClickHouse JDBC的sink建议开启jdbc_use_compression减少跨网络的数据体积。实测100MB的文本数据开压缩后传输量降到20MB左右对延迟敏感的场景很明显。3.6 性能调优与常见坑我用真实数据踩出来的结论这部分是我最想写、也是干货密度最高的部分。因为在同步链路里“看着一切正常”和“真的高性能稳定运行”之间隔着一堆反直觉的细节。坑一Binlog里的大事务会拉高同步延迟。MySQL一个事务如果涉及几十万行更新Flink CDC会等到事务提交后再一次性发给下游。这是Binlog本身的机制事务内的变更在提交前不可见。解决思路是把大事务拆小从业务源头改。比如批量更新脚本改成分批提交每批5000行。如果业务改不了就在Flink端对这个大事务的事件做缓冲当它提交后再下发。注意这种方式对读取端的吞吐冲击很大建议单独监控“source currentFetchEventTimeLag”指标。坑二ClickHouse写入端的背压表现不像Kafka那样平滑。Flink的sink在阻塞时会不断累积反压如果ClickHouse端merge不及时too many parts错误会频繁触发。遇到这种情况常规操作是把sink的并发度和batch size调大而不是去调ClickHouse的max_partitions_per_insert_block。更彻底的做法是给目标表设置一个“先落临时表再定期把临时表分区替换到正式表”的批处理链路。虽然多了一步但对大吞吐同步场景的稳定性有质的帮助。坑三时区问题。这是最常见的脏数据来源。MySQL的datetime不带时区timestamp带时区。当Flink JDBC读取MySQL时默认用JVM时区而ClickHouse端写入DateTime时按服务器时区来解释。如果两端时区不一致同步过去的时间会偏移8小时。解决方式在Flink CDC的JDBC连接参数中加上serverTimezoneAsia/Shanghai在ClickHouse JDBC连接参数中设置use_time_zoneAsia/Shanghai并且保证ClickHouse服务端配置了正确的timezone。坑四表结构里如果有ENUM或SET类型同步到ClickHouse会变成什么老版本的JDBC映射可能把它们映射成String但Flink CDC在捕获MySQL DDL时会把它们解析成String所以我建议建表时统一在目标端改成String或LowCardinality(String)。后者对值的数量有限制适合状态明确的字段能大幅提升过滤和聚合速度。坑五不要在主键或排序键中使用可空字段。这是ClickHouse的硬性规则。MySQL表的主键不会为NULL但如果有字段参与自定义排序键而这些字段在业务上允许NULL那么Flink端写入时要提前转换。我的做法是在Flink SQL里写COALESCE(field, )或COALESCE(field, 0)让排序键字段永远非空。下面是整个同步任务的Flink SQL骨架读者可以直接参考修改-- 创建源表使用Flink CDC连接MySQL CREATE TABLE mysql_orders ( id BIGINT, user_id BIGINT, order_status STRING, total_amount DECIMAL(10, 2), create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 10.10.1.20, port 3306, username flink_cdc_user, password xxxx, database-name trade, table-name orders, scan.startup.mode initial, debezium.snapshot.fetch.size 4096 ); -- 创建ClickHouse目标表 CREATE TABLE ch_orders ( id BIGINT, user_id BIGINT, order_status STRING, total_amount DECIMAL(10, 2), create_time TIMESTAMP(3), update_time TIMESTAMP(3) ) WITH ( connector jdbc, url jdbc:clickhouse://10.10.1.30:8123/trade, table-name orders, username default, password xxxx, sink.buffer-flush.max-rows 2000, sink.buffer-flush.interval 3s, jdbc.connection.max-retry-timeout 300s ); -- 流式写入 INSERT INTO ch_orders SELECT id, user_id, order_status, total_amount, create_time FROM mysql_orders;这段SQL里scan.startup.modeinitial表示启动时先做全量快照再进入增量这是全量增量一体的关键配置。如果只想要增量可以改成latest-offset。从实际运维角度我还建议在Flink UI上重点盯两个指标currentFetchEventTimeLag数据从MySQL产生到Flink捕获的延迟和currentEmitEventTimeLag捕获后到实际写出的延迟。正常情况下百毫秒级算健康如果持续秒级先检查MySQL端是否有大事务再查ClickHouse端写入是否有背压。另外对于ClickHouse端经常出现的Too many parts错误除了调大batch size还有个兜底操作把目标表的merge_with_ttl_timeout调大一点给后台merge多一点时间。这个方法治标但能保证链路不断。4. 新版本对同步场景的额外影响为什么25.12更值得升级单独把25.12和Flink CDC这套方案的兼容性拎出来说是因为新版Pipeline模型对高吞吐写入的调度确实有变化很多人升级后会观察到“写入变快”或“Part数量变少”的表象但背后的机制值得理解。Pipeline模型下插入和merge任务共用同一套线程池调度。相比旧模型每个查询独占几条线程Pipeline模型能让空闲查询让出CPU给写入和merge任务。所以在Flink持续大批量写入的场景里25.12的并发写入调度更积极。我做过一次小幅基准测试同样设置20个并发写入任务每个任务批量插入2000条25.12和24.8相比写入P99延迟下降约30%而且CPU峰值没有上升。这个差异在纯单条写入场景里不明显但在Flink这种“批量轰炸”模式下感受很清晰。零拷贝复制对同步链路的影响则体现在另一个层面。假设我们搭建了一套“Flink实时写入一个业务副本另一个副本用于BI查询”的架构在老版本里两个副本需要各自落盘数据文件占双倍磁盘。25.12的零拷贝复制在三个节点上做测试同位副本的数据文件共用磁盘占用减半副本间的同步延迟也缩短到秒级以内。这对那些“同一份数据既要给实时报表查询又要给离线分析跑批”的团队非常友好。不过零拷贝复制有个硬前提ClickHouse必须运行在支持copy_file_range的文件系统上。常见对象存储S3、OSS在实现上对这类系统调用支持不完整所以如果集群的数据目录在对象存储挂载盘上零拷贝复制可能需要单独验证或关闭。我的建议是先把零拷贝复制用在一组新构建的物理盘或高效云盘集群上跑一周观察Part移动情况和系统日志。稳妥起见线上核心集群可以等到25.12.x的小版本修复几轮后再全量铺开。5. 一次完整的升级与同步联调过程从准备到验证的实操记录理论说再多不如把一次完整的操作过程记录下来。下面是我在测试环境里做的升级同步联调供读者复现时对照。步骤一准备测试环境。我搭建了一个三节点ClickHouse集群版本从24.3升级到25.12。系统是Ubuntu 22.04数据目录用NVMe云盘。Flink Standalone集群部署在一台8C16G的机器上Flink CDC安装在$FLINK_HOME/lib目录。步骤二备份与降级预案。升级前用clickhouse-client执行BACKUP TABLE xxx TO Disk(backup_disk, path)做全量备份。注意ClickHouse的BACKUP命令在25.12里支持增量备份但对于升级场景全量备份更保底。同时保留上一版本的二进制包万一升级后发现严重问题可以直接替换二进制回滚。步骤三升级ClickHouse。我用的Tarball方式升级解压新版二进制后把旧版的config.xml、users.xml保留。启动新版时system.warnings表会列出一些“配置项已被移除”的提示比如我在旧配置里写了allow_experimental_lightweight_update1/enable_experimental_lightweight_update新版直接提示不再需要。清理掉这些废弃配置后集群正常启动。步骤四验证新特性。新建一张测试表插入100万行数据开启Pipeline模型的默认设置跑一个带窗口函数的查询SELECT user_id, row_number() OVER (PARTITION BY user_id ORDER BY create_time DESC) AS rn FROM orders_test LIMIT 1000;查询顺利返回结果。对比24.3版本执行时间从大概2.4秒降到1.2秒。步骤五部署Flink CDC同步任务。按第三节里的SQL骨架在Flink SQL Client中逐段执行建表语句。启动同步前先确认MySQL的Binlog格式是ROW且binlog_row_imageFULL。然后执行INSERT INTO ch_orders SELECT ...立刻触发全量同步同步完成后自动进入增量。步骤六数据校验。通过对比MySQL和ClickHouse的记录数、SUM聚合值来校验。我用以下SQL做全量校验-- MySQL侧 SELECT count(*), sum(total_amount) FROM trade.orders; -- ClickHouse侧 SELECT count(*), sum(total_amount) FROM trade.orders;两个值在同步完成且merge结束后差值应接近0。由于实时同步下数据仍在变化我选择在低峰期且暂停写入业务的窗口做校验或者用max(update_time)做增量比较。步骤七验证容错。手动杀掉Flink TaskManager进程观察JobManager是否按Checkpoint恢复任务。我在测试中故意Kill一个TaskManager大概20秒恢复后通过system.query_log检查ClickHouse端是否出现重复写入。由于目标表用的是MergeTree且没有幂等控制重复写入确实出现了几条。随后我改用ReplacingMergeTree表并重新同步重复问题消除。6. 调优策略和踩坑边界适合上线前再回头看的内容最后这部分写给准备上生产的人。和ClickHouse 25.12与Flink同步链路相关我在近一个月测试中总结了下面的调优优先级预热分区与分区键分区的选择直接影响写入和查询的性能。对于时间字段如果查询多集中于最近7天建议分区按天设置然后配合TTL清理旧分区。不要为了“省事”把所有数据塞进一个分区否则单分区Part数量可能几千个查询和merge都会变得缓慢。调大sink的batch size这是最立竿见影的优化。Flink JDBC Sink默认buffer-flush.max-rows是100我通常调到2000到5000。前提是ClickHouse端能接受这个批次的写入延迟否则反而会导致反压。实测2000是一个比较稳的中间值。开启SSL与压缩相比明文传输不用纠结性能损耗。在现代CPU上压缩和解压的开销通常小于网络IO的节省。ClickHouse JDBC的sslon和jdbc_use_compressiontrue建议直接开。监控链路指标Flink UI中的busyTimePerSecond和outputQueueLength是观察sink是否瓶颈的最直观指标。对ClickHouse端则建议关注system.metric里的BackgroundPoolTask和MemoryTracking。再提几个容易产生“我上了生产才后悔”的边界场景。一是数据回填如果同步任务暂停了很久恢复后Flink会从最近一次Checkpoint继续但期间MySQL积累的大量Binlog变更会一次性流入ClickHouse可能出现短时“挤兑”。处理办法是暂停时同时停掉业务写入或者接受延迟恢复不要频繁启停。二是MySQL侧大事务这个问题前面提过它会让Flink端source延迟突增但不会丢数据关键是监控出来之后在业务侧推动拆分。三是ClickHouse侧磁盘慢IO如果数据目录的IOPS不够Flink写入再努力也会被ClickHouse端背压拖住这时升级磁盘比调Flink参数有效得多。最后建议所有团队在上线前准备一个“数据比对工具”用Python或Shell脚本定期对比MySQL和ClickHouse的计数、SUM、MAX、MIN这比在Flink端做任何复杂的断点校验都更能发现问题。把同步链路的稳定建立在“日常自动校验警报”之上比追求完美的Exactly-Once更实际。
返回列表