恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
ClickHouse实时洞察:Flink同步MySQL与部署选型实践
首页
资讯中心
/
ClickHouse实时洞察:Flink同步MySQL与部署选型实践
ClickHouse实时洞察:Flink同步MySQL与部署选型实践
发布时间:2026/10/6 3:02:12
ClickHouse 这名字在我这几年的大数据项目里出现频率越来越高。最早接触它是接手一个数据平台改造几十亿行用户行为数据业务方要实时汇总出核心指标MySQL 加 Redis 的组合已经撑不住离线 Hive 又只能回答昨天的问题。当时我对比了几种方案最后把 ClickHouse 推到线上效果比预想中好不少。这篇文章围绕“大数据实时洞察”这条主线把我自己踩过的坑和验证过的方法整理一遍重点是 Flink 实时同步 MySQL 数据到 ClickHouse 的链路、Linux 上部署 21.8 版本的注意点、跟 Doris 的选型对比以及数据大屏这类下游应用的落地思路。想在团队里接入实时分析能力的同学应该能找到几张能直接抄的作业。1. 实时洞察的本质先把“查得慢”这件事解决掉1.1 为什么 OLTP 数据库在聚合分析面前会露怯很多业务系统一开始跑得挺好订单表每天也就几万行一个GROUP BY查询再慢也就几百毫秒。但当一张表积累到几亿行同样的统计查询可能要几十秒甚至几分钟。这不是 SQL 写得不好而是存储引擎的设计目标不一样。MySQL 这类 OLTP 数据库基于 B 树和行式存储它的核心优势是单行点查、高并发事务、强一致。做聚合统计的时候它不得不把命中的整行数据读出来哪怕你只需要其中两列。而且GROUP BY、排序、临时表的创建都会消耗大量内存和磁盘 IO行数一上来瓶颈立刻出现。更麻烦的是OLTP 数据库的索引主要服务点查和范围过滤对聚合统计基本帮不上忙——GROUP BY除了在索引覆盖特殊场景外大多数时候是绕不开全表扫描的。我见过最多的做法是先上缓存不行就建汇总表还不行就加索引最后被业务催着上 Hadoop。每一个方案都能缓解问题但都没有真正解决“海量明细秒级聚合”的底层矛盾。1.2 列式存储和向量化执行ClickHouse 的快不是玄学ClickHouse 解决这个问题的思路很直接既然分析查询只需要少数几列那就按列存。按列物理存储之后查询SUM(pay_amount)时只读取这一列的压缩数据而不是把整行数据拖出来。配合 LZ4、ZSTD 这类压缩算法磁盘 IO 可以下降一个量级很多场景下甚至是几个量级。更关键的是向量化执行。ClickHouse 在内存里把连续的一批数据交给 CPU 处理通过 SIMD 指令同时对多条数据执行同一操作。相当于把“扫描 聚合”这个过程变成了流水线作业而不是逐行解释执行。实战中同样的聚合查询MySQL 跑 20 秒的ClickHouse 经常在几十毫秒到几百毫秒内返回这种差距不需要复杂的理论解释就是执行模型本身的差异。1.3 MergeTree 家族实时洞察场景里的表引擎选型ClickHouse 的快一半靠列式存储另一半靠表引擎。MergeTree是基础引擎数据按分区和排序键组织ReplacingMergeTree支持按版本去重适合处理 MySQL 同步过来的更新数据SummingMergeTree和AggregatingMergeTree可以在后台合并时做预聚合适合指标汇总场景。以订单同步为例我通常会这样建表CREATE TABLE ods_order ( order_id UInt64, order_status String, province String, city String, pay_amount Decimal(18, 2), order_time DateTime64(3, Asia/Shanghai), update_time DateTime64(3, Asia/Shanghai) ) ENGINE ReplacingMergeTree(update_time) PARTITION BY toYYYYMMDD(order_time) ORDER BY (order_id);ReplacingMergeTree(update_time)表示合并数据时同一个order_id保留update_time最大的那行。分区字段用order_time是为了让查询在按天裁剪分区时能把 IO 降到最低。排序键选业务主键是为了让同一条订单的多次变更尽量落在相邻位置合并去重更高效。想清楚这些再去搭数据链路就不会盲目。实时洞察的根基是查询模型和存储模型先选对否则后面加再多缓存都是治标不治本。2. Flink 消费 MySQL Binlog同步到ClickHouse一条实时链路的完整拆解2.1 为什么选 Flink CDC 而不是 Canal 全家桶实时洞察的前提是有实时数据。大多数业务系统的数据源还是 MySQL比如订单库、用户库。传统做法是每天凌晨跑批同步但业务方要的是分钟级甚至秒级延迟。同步工具可选的有几类。Canal 负责订阅 MySQL binlog把变更消息投递到 Kafka 或 MQ但下游还需要自己写一个消费程序去应用这些变更。Debezium 本质上是做同样的事情只是套在 Kafka Connect 生态里。Flink CDC 的好处在于binlog 捕获、状态管理、checkpoint、断点续传都封装好了而且可以直接用 Flink SQL 定义同步任务把“同步 计算 写入”放在同一条链路上。如果团队本身已经有 Flink 平台我会优先选 Flink CDC。如果只是很轻量的数据搬运追求极简运维Canal 加一个自研同步器也可以。但要注意一旦涉及复杂的字段转换、多表 join、迟到数据处理Flink 的优势会明显放大。这里有一个容易被忽略的点Flink CDC 的增量快照。任务启动时先扫描全量数据再无缝切到 binlog 增量整个过程对源库的影响比较小不会像传统SELECT *那样长时间锁表。底层逻辑是分片枚举加无锁读取靠记录 binlog 位置来保证断点续传。对生产环境来说这一项能力比什么都重要。2.2 一条可直接参考的 Flink SQL 同步方案用 Flink SQL 写同步任务代码量很小。先定义 MySQL CDC 来源表CREATE TABLE mysql_order ( order_id BIGINT, order_status STRING, province STRING, city STRING, pay_amount DECIMAL(18, 2), order_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 192.168.1.10, port 3306, username bigdata, password xxx, database-name trade, table-name order_info, scan.startup.mode latest-offset, debezium.snapshot.mode initial );再定义 ClickHouse 目标表。这里要注意Flink SQL 原生 JDBC Connector 对 ClickHouse 的支持并不算特别好生产环境我更建议自己写一个 DataStream Sink用clickhouse-jdbc的批量写入能力或者直接用社区维护的clickhouse-flinkConnector。在 SQL 里示意的话大概长这样CREATE TABLE ck_order ( order_id BIGINT, order_status STRING, province STRING, city STRING, pay_amount DECIMAL(18, 2), order_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:clickhouse://192.168.1.20:8123/trade, table-name ods_order, sink.buffer-flush.max-rows 5000, sink.buffer-flush.interval 10s ); INSERT INTO ck_order SELECT * FROM mysql_order;如果自己写 Sink有四个要点必须盯住攒批单批尽量达到 5000 行以上或者攒 10 秒再写不要一条一条插。批量失败重试重试两次仍失败就把数据写到死信表避免阻塞整条链路。并发度控制写入 ClickHouse 的并发建议控制在 4 到 8 之间并发太高会产生大量小分区文件后台 merge 跟不上磁盘和查询都会受影响。Checkpoint开启 Flink CheckpointClickHouse 本身没有事务性写入最终一致性靠目标表的幂等设计来兜底。2.3 类型映射与数据一致性校验MySQL 和 ClickHouse 字段类型不是一一对应的同步前最好拉一张映射表MySQL 类型ClickHouse 类型说明TINYINTInt8 / UInt8布尔字段可映射为 UInt8INTInt32常规整数BIGINTInt64常用主键、时间戳VARCHAR / CHAR / TEXTString无需关心长度DATETIME / TIMESTAMPDateTime64(3)建议统一毫秒精度DATEDate32存储日期DECIMAL(p, s)Decimal(p, s)精度要一致这里容易踩的坑是时区。MySQL 的DATETIME不带时区信息ClickHouse 的DateTime64又是带时区语义的。如果服务器时区不一致查出来的时间可能整体偏移 8 个小时。我的做法很粗暴同步链路里所有时间字段统一转成Asia/Shanghai时区在 Flink SQL 里显式转换不依赖服务器默认时区。数据一致性校验也不能省。我会建一张sync_check表每条同步任务定期记录源表和目标表的COUNT、SUM、MAX、MIN每 5 分钟跑一次对比。偏差超过阈值就触发告警再抽样明细定位。很多同步问题不是第一天才爆发的而是跑了一周之后字段枚举值变化或者某条数据格式异常才暴露没有校验机制就是在裸奔。2.4 实时同步中的版本冲突与幂等处理MySQL 会有UPDATE和DELETE但 ClickHouse 默认不擅长行级更新。最稳的做法就是目标表用ReplacingMergeTree在同步 SQL 里把update_time作为版本字段。合并之后同一个业务主键只保留最新版本。至于删除事件我在生产环境一般不做物理删除而是把is_deleted字段置为 1查询时统一过滤。这样既躲开了 ClickHouse 高频删除的短板也保留了审计历史。真需要物理删除的场景可以用 21.8 引入的 Lightweight Delete或者在表结构上用CollapsingMergeTree做软删除。另一个细节是主键复用。有些业务系统的主键不是严格的“创建后不可变”比如订单号回滚之后会被复用。这种情况如果只靠一个字段去重新老订单互相覆盖就会出现脏数据。我的做法是除了业务主键再保留一个自增代理键或者引入version字段确保排序键的唯一性由“业务主键 版本号”共同保证。3. Linux部署ClickHouse 21.8.15.7从单机到分片集群的落地细节3.1 为什么我会锚定 21.8 这个版本ClickHouse 的版本迭代非常快几乎每个月都有新版本。生产环境我最看重的是 LTS 版本。21.8 是官方定义的长期支持版本之一在后续很长一段时间内依然能收到修复补丁同时引入了 Lightweight Delete、Projection、窗口函数优化等实用特性。很多网上的教程是拿当时最新版写的配置路径、语法细节跟 21.8 差别不小。如果照着搜来的文章配很容易踩到“参数不存在”“语法不被支持”这种问题。我建议生产环境认准一个 LTS 版本部署和排障都以官方文档为准不要盲目追新。3.2 单机部署的系统参数和关键配置以官方 tgz 包为例安装步骤不复杂tar -xzf clickhouse-common-static-21.8.15.7.tgz cd clickhouse-common-static-21.8.15.7 ./install/doinst.sh tar -xzf clickhouse-server-21.8.15.7.tgz cd clickhouse-server-21.8.15.7 ./install/doinst.sh tar -xzf clickhouse-client-21.8.15.7.tgz cd clickhouse-client-21.8.15.7 ./install/doinst.sh真正决定稳定的不是安装而是系统参数和配置。有三个东西我会在部署时第一时间设置/etc/security/limits.conf把nofile调到 262144nproc调到 65535否则高并发查询很容易报句柄不足。vm.max_map_count设置到 262144。ClickHouse 在内存映射和文件映射上很激进默认值太小会出现Cannot mmap错误。数据盘独立数据目录放到单独的 SSD 上别跟系统盘混在一起。ClickHouse 对磁盘 IO 非常敏感机械盘和系统盘混用会让 merge 和查询互相拖累。config.xml里重点看这几项listen_host0.0.0.0/listen_host path/data/clickhouse//path tmp_path/data/clickhouse/tmp//tmp_path max_server_memory_usage物理内存的60%/max_server_memory_usage mark_cache_size1073741824/mark_cache_sizemax_server_memory_usage建议给到物理内存的 50% 到 70%剩下的留给操作系统和文件缓存。这个参数如果配得太高碰到几个并发大查询很容易直接把内存打满触发 OOM。users.xml里默认用户default一定要改密码别让默认空密码上线users default password_sha256_hex这里填SHA256哈希/password_sha256_hex networksip::/0/ip/networks profiledefault/profile quotadefault/quota /default /users3.3 多副本与分片什么时候开始上集群很多人一开始就想上集群其实多数场景下单个节点已经能扛住很大的查询压力。我的判断标准很简单单机数据量在几 TB 以内查询并发没有爆炸式增长就先单节点跑。上副本是为了高可用不是为了性能。ClickHouse 单节点本身的写入和查询性能已经足够强加副本只会增加 ZooKeeper / Keeper 的运维成本。决定上副本之后用ReplicatedMergeTree就对了。它需要外部协调服务21.8 已经支持 ClickHouse Keeper可以替代 ZooKeeper。说实话Keeper 的部署和维护成本比 ZK 低不少新环境我一般直接用它。每个节点上配置好宏macros shard01/shard replicareplica_1/replica /macros再通过ON CLUSTER的方式在建表语句里统一创建CREATE TABLE shard_order ON CLUSTER ck_cluster ( order_id UInt64, order_status String, pay_amount Decimal(18, 2), order_time DateTime64(3, Asia/Shanghai) ) ENGINE ReplicatedMergeTree(/clickhouse/tables/{shard}/order, {replica}) PARTITION BY toYYYYMMDD(order_time) ORDER BY order_id;接着建一张分布式表CREATE TABLE all_order ON CLUSTER ck_cluster ENGINE Distributed(ck_cluster, default, shard_order, rand());分片键的选择很关键。订单号散列和用户 ID 取模是常见选择目的都是让数据尽量均匀地落到不同分片。最怕的是拿时间字段做分片键业务高峰时段所有数据都写进同一个分片查询时某个节点被打满其他节点在摸鱼。3.4 行级权限和列级权限的开源实践大数据平台做到一定规模权限治理就是必答题。好在 ClickHouse 内置的 RBAC 已经能做不少事。行级权限用 row policy 实现。举个例子销售运营只能看自己区域的订单CREATE ROW POLICY pol_sales ON trade.orders FOR SELECT USING sales_region currentUser();列级权限则通过授权语句控制。给张三只开放订单号和金额不让他看用户手机号GRANT SELECT(order_id, pay_amount) ON trade.orders TO zhangsan;这套内置方案在数据量可控、角色清晰的场景下完全够用。如果团队希望做更集中的权限管理社区里有 Apache Ranger 的 ClickHouse 插件但成熟度不如 Hive 那边的生态。我的落地建议是内部报表场景直接用内置 RBAC 加 LDAP 认证如果是 SaaS 多租户场景在数据接入阶段就把租户 ID 作为过滤条件写入 row policy 或视图而不是让每一条查询请求都到应用层做一遍权限判断。权限只放在应用层的后果是BI 工具、临时查询、实时任务来源一多一定会漏。3.5 日常监控和故障自救ClickHouse 自带的system表信息量非常大。我最常用的有三张system.query_log看慢查询和内存消耗system.metrics看当前运行状态system.events看各类事件计数。排查慢查询时直接跑SELECT query, query_duration_ms, read_rows, memory_usage FROM system.query_log WHERE event_time now() - INTERVAL 1 HOUR ORDER BY query_duration_ms DESC LIMIT 20;监控层面用 ClickHouse-Exporter 加 Prometheus 加 Grafana重点盯几个指标CPU、磁盘使用率、merge 队列长度、内存占用。我碰到过的绝大多数线上故障在 merge 队列变长和内存曲线异常时就已经有征兆了。别等业务方来投诉自己先把告警配好。4. 选型不迷路ClickHouse和Doris的对比与实践感受4.1 两者到底哪里不一样ClickHouse 和 Doris 都是列式存储、MPP 架构、面向 OLAP 的分析引擎也都能承载实时数仓场景。但它们的定位差异很大放到同一张表格里看会更清楚对比维度ClickHouseDoris数据模型MergeTree 家族分区 排序键Duplicate / Aggregate / Unique 三种模型SQL 兼容方言较独特高度兼容 MySQL 协议更新能力靠 ReplacingMergeTree / Collapsing 家族Unique Key 模型更成熟JOIN 复杂度适合宽表和轻量 join复杂 join 更稳下游生态外部表、物化视图丰富管理界面更完善BI 工具接入简单运维组件clickhouse-server KeeperFE、BE、Broker 多组件典型场景海量日志、行为分析、实时大屏实时报表、数据服务、多表查询ClickHouse 最大的优势是单表聚合和点查速度极快尤其在超大宽表、超多分区这类场景下表现非常稳定。Doris 的强项在“接近 MySQL 的使用体验”BI 工具连它就像连 MySQL多表 join 的查询规划器更完善Unique Key 模型对高频更新友好得多。4.2 我在 ClickHouse 上踩过的坑第一坑是 JOIN。两个大表做 join默认的 hash join 算法会吃掉大量内存配置不当直接 OOM。后来我把可下推的维度都做成低基数字典能不用 join 就不 join尽量用宽表替代。ClickHouse 的风格是想办法把数据模型设计成“单表可查”而不是依赖 join 解决问题。第二坑是 DELETE 和 UPDATE 的误操作。21.8 虽然支持 Lightweight Delete但大批量删除会触发 mutation 队列短时间可能阻塞查询。所以同步链路里的逻辑删除尽量别用物理删除。第三坑是分区粒度。给订单表按天分区没问题但如果业务只查最近 30 天的数据按天分区让查询裁剪很高效反之一天一个分区但每天只有几千行part 数量会失控merge 永远跑不完。4.3 什么场景下我更倾向 Doris如果项目的核心需求是“一堆 BI 报表 多表关联 运维界面友好”Doris 的优势是实打实的。MySQL 协议带来的兼容性让 Tableau、帆软这类工具几乎零成本接入Unique Key 模型配合 Stream Load 或 Routine Load从 Kafka 导入数据后天然支持主键更新。对于要从零建设一套实时数仓、面对的是几十张明细表和大量业务报表的团队Doris 的学习曲线更友好。反过来如果数据规模动不动就到 PB 级查询以日志分析和行为分析为主需要极致压缩比和并发控制能力ClickHouse 更合适。而且它和 Hive 生态的互补性很好离线数仓已经用 Hive 的团队加一个 ClickHouse 做实时加速几乎不用改变原有架构。4.4 我的选型建议给一个可以照抄的判断逻辑已有 Hive 加 Spark 离线数仓只是想给实时洞察提速选 ClickHouse。从零建设报表查询多、join 复杂、团队前期人手不足选 Doris。想用最小的组件栈快速验证先上 ClickHouse 单节点跑真实业务查询跑通再考虑副本。集团有统一的数据权限、数据服务规范Doris 的环境更好融入。我最想说的是不要同时上两套 OLAP 引擎。数据团队资源有限维护两套引擎的成本会摊薄做业务优化的精力。选一个拿真实数据量跑一个月 POC再做决定。5. 从秒级结果到可读洞察实时大屏与下游应用落地5.1 大屏不是炫技它让实时洞察有了出口实时洞察最终要落到业务动作上数据大屏是一个很直观的出口。以网约车项目为例平台要实时看订单总量、活跃司机数、区域热力、各时段成交金额。这些指标反映的是“当前正在发生什么”大屏把 ClickHouse 处理后秒级返回的结果转成运营和管理人员看得懂的图表。如果背后还是离线 T1 数据大屏就只是摆了一个好看的壳。5.2 一张实时大屏的数据链路设计我常用的链路是ClickHouse 负责聚合计算后端服务定时查询Redis 做结果缓存前端用 ECharts 做轮询展示。前端每 5 秒刷新一次请求先打到后端接口。后端优先读 Redis缓存过期或不存在时再查 ClickHouse并把结果写回缓存缓存时间控制在 30 秒左右。为什么要加这一层因为大屏往往不止一块多块大屏同时轮询时不经过缓存会把 ClickHouse 打成高频重复查询。ClickHouse 擅长的是重分析 SQL不是每秒几十次的短查询。加一层 Redis既保证刷新体验又限制了下游压力。指标 SQL 长这样SELECT toStartOfMinute(order_time) AS minute, count() AS order_cnt, sum(pay_amount) AS gmv FROM ods_order WHERE order_time now() - INTERVAL 30 MINUTE GROUP BY minute ORDER BY minute;这条 SQL 返回最近 30 个分钟点的订单量和成交金额前端直接画折线图。区域热力图则把GROUP BY换成城市加五分钟窗口逻辑完全一样。5.3 实时大屏的优化思路第一件事是预聚合。实时大屏查询的粒度往往固定比如分钟、小时、城市。与其每次都扫明细表不如建一张分钟级预聚合表由物化视图或定时任务从明细表聚合写入。查询走预聚合表响应时间能降一个数量级明细表也不会因为大屏查询增加额外压力。第二件事是巧妙利用 Projection。21.8 已经支持 Projection相当于在原始表内部预存了一份针对常见聚合的投影数据。查询时如果匹配到投影就不需要重新扫描和聚合自动减少计算成本。这个功能对“多张大屏 固定查询模式”的场景非常实用。第三件事是把计算尽量下推到 SQL。比如同比环比与其把几万行明细拉到 API 层算不如直接用窗口函数在 SQL 里算好API 层只是透传结果。既省流量又省代码量。第四件事是给大屏查询设置边界。大屏类 SQL 一定要有明确的时间范围条件比如只看最近 30 分钟防止因为前端某次参数异常误触全表扫描。我还会在查询里设置max_execution_time超过阈值直接中止宁可这次刷新失败也不能让一个大屏查询拖垮整个集群。5.4 实时洞察的另一种出口告警和数据服务大屏是给人看的告警是给系统用的。实时洞察完全可以做成自动化的异常检测每 10 秒跑一次规则 SQL比如订单量同比骤降、支付成功率低于阈值触发后通过 Webhook 推给值班群。这类查询的特点是次数多、单次轻最好在 ClickHouse 用户配置里单独限制内存和超时避免被某一次异常数据打爆。数据服务化是更进一步的方向。把 ClickHouse 的指标查询封装成标准 API给客服、运营、财务等部门使用。鉴权直接复用 row policy不同角色天然只能查到有权限的数据。实时大屏只是“洞察”的一个前端真正让这套体系稳定跑起来的是数据延迟监控、异常标记和结果留痕这些基本功。最后分享一些我的实际感受。实时洞察这件事最难的不是选引擎而是把数据质量守住。同步链路跑三天就会碰到字段类型对不上、时区错乱、偶发重复数据先把这些解决掉比换更高版本的引擎有价值得多。ClickHouse 在我眼中更像一个“查询加速器”而不是万金油。单节点跑 POC把真实业务查询都过一遍再决定要不要上副本这条路我走过不止一次至今没后悔过。