ARTICLE DETAIL

资讯详情

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

Flink + ClickHouse 亿级实时数据分析平台:部署、同步与调优实践

Flink + ClickHouse 亿级实时数据分析平台:部署、同步与调优实践 简介这是一份基于Flink与ClickHouse构建的亿级电商实时数据分析平台完整项目覆盖PC端、移动端与小程序三端应用面向大数据方向的学生、开发者及毕业设计使用者。包内含完整前后端源码、部署文档、配置说明及辅助资料共1136个文件以Java、Vue、JavaScript等代码文件为主辅以HTML页面、CSS样式、PNG图片、XML/JSON配置及少量PSD设计源文件压缩包整体仅7.07MB目录结构层次分明便于按模块检索、运行与二次开发。项目涵盖实时数据采集、ETL处理、指标体系计算、可视化大屏展示等典型环节代码经过测试运行验证并已通过导师指导与答辩评审评分95分适合作为毕业设计、课程设计、项目初期立项演示也可作为学习Flink实时计算与ClickHouse分析引擎的完整实战案例。目前已有94人学习下载能够帮助开发者快速理解电商实时数仓从数据接入到可视化展示的工程化落地方式。1. 为什么偏偏是 Flink ClickHouse这个组合到底解决什么问题做电商实时数据分析平台最难的不是凑一套报表而是数据从 MySQL、埋点日志、小程序端进来到最终能在看板里秒级看到 GMV、订单量、转化漏斗这中间每一步都可能成为瓶颈。这个项目标题点名的 Flink ClickHouse恰好是目前做亿级数据实时分析最稳的一套搭配Flink 负责把源源不断的订单、点击、支付事件接进来并做实时计算ClickHouse 负责把结果数据存成可秒查的列式存储。它们解决的是同一个问题——数据量一大传统关系库扛不住、离线数仓又不够快而你需要在秒级完成从数据产生到分析结果呈现的完整链路。这套方案适合正在做电商数据中台、实时大屏、用户行为分析的开发团队也适合准备把 Flink 和 ClickHouse 作为简历核心技能的从业者。接下来我按部署、数据同步、调参与排错的顺序把这套方案完整拆开。2. 先看懂技术选型Flink 和 ClickHouse 在实时链路里各自扮演什么角色2.1 Flink 不只是“快”而是把无界数据流变成了可计算的编程模型很多人在接触 Flink 时有一个误解觉得它只是比 Spark Streaming 吞吐量高一些的流处理框架。实际上 Flink 真正厉害的地方在于它重新定义了“流”和“批”的关系所有数据都可以被当作无界流来处理一个订单进来、一个点击发生都是一条数据流事件。它通过 checkpoint 机制把运行状态定期打到持久化存储里当某个 TaskManager 挂了就会自动从最近一次 checkpoint 恢复做到精确一次的处理语义。这就是为什么在电商场景里一个支付成功事件如果被重复计算GMV 就会虚高而 Flink 可以做到事件不丢、不重、不乱序。在这个项目里Flink 承担的不只是聚合计算。它还要做清洗把前端埋点里的时间戳统一格式把缺失的用户 id 过滤掉把重复的取消订单事件剔除然后把结果输出到下游。这个过程用 Flink 的 DataStream API 或者 Table API 都能完成。我用得比较多的是 Table API因为电商实时看板里其实大部分需求都是 group by 加窗口用 SQL 写更直观后面想加一个维度也更好改。要注意的是Flink 的窗口并不是等数据全部到齐才开始计算它是按照事件时间或者处理时间去触发。电商场景里必须用事件时间也就是订单里真实的发生时间否则客户端网络抖动导致的数据延迟到达会把统计结果搅乱。2.2 ClickHouse 的列式存储和 MergeTree 家族为什么适合亿级聚合结果ClickHouse 能在亿级数据上做到秒级返回靠的是两个核心设计列式存储和向量化执行。列式存储意味着查询一列聚合值时不需要像 MySQL 那样把整行数据读入内存而是只读取参与计算的列数据磁盘 IO 大幅度下降。向量化执行则是把原来一条一条处理数据的循环改成了一次处理一批数据让 CPU 的 SIMD 指令集充分发挥作用。这两点叠加使得 ClickHouse 在处理 group by、count distinct 这类分析型查询时性能远超传统关系型数据库。在这个实时分析平台里ClickHouse 的存储引擎选择也有讲究。 ReplacingMergeTree 适用于 upsert 场景也就是同一主键会有多条记录、只保留最后一条的情况。SummingMergeTree 则适合按维度预先聚合的累计表Flink 写入明细后ClickHouse 后台会自动把相同维度键的数据合并求和。电商场景里最常用的是这两类引擎的组合实时大屏查询用 SummingMergeTree 预聚合表数据稽查和订单明细查询用 MergeTree 原始表。分区键一般按天设置亿级数据量下按天分区既能保证查询裁剪掉无关数据又能让 TTL 过期清理变得非常方便。2.3 为什么这个场景不适合直接用 Doris 或者 Elasticsearch搜索热词里大量出现 doris 和 clickhouse 的选型我确实在实际项目中踩过这个弯。Doris 的强项在于它把 FE 和 BE 分离支持高并发查询并且在 join 能力上比 ClickHouse 友好适合相对固定的报表系统。但电商实时分析的一大特点是写入峰值非常陡比如大促零点那几分钟Flink 计算结果快速往库里灌JDBC 批量写入一旦超过 ClickHouse 的 merge 节奏查询侧会出现短暂抖动ClickHouse 对这种高吞吐写入的容忍度比 Doris 更宽。Elasticsearch 则更适合搜索和全文匹配用它做亿级聚合会把聚合节点压到内存崩溃而且 ES 的存储成本大概是 ClickHouse 的三到五倍量一大光磁盘开销就受不了。所以这套实时分析平台选 ClickHouse 的核心理由是写入吞吐高、聚合查询快、存储成本可控部署包 21.8 LTS 版本也比较稳后面所有部署细节都基于这个版本展开。3. 在 Linux 上从零部署这套实时分析底座Flink、Kafka 与 ClickHouse 的拓扑与参数3.1 先画一张最小可用部署拓扑再动手下载安装我一般不会一上来就搭三台机器的集群而是先单机把链路跑通再扩展成集群。单机版部署拓扑是这样的一台 16 核 32G 的 Linux 服务器上面跑一个 Flink Standalone 集群一个 JobManager、两个 TaskManager 进程、一个 Kafka 单节点、一个 ClickHouse 单机实例。这个配置跑亿级数据的离线回放可能有点挤但用来验证实时链路、压测 JDBC 写入、调试部署文档里给的 SQL完全够用。生产环境建议至少三台一台放 JobManager 和 ClickHouse两台放 TaskManager 和 Kafka因为 Flink 的 TaskManager 吃 CPU 最凶Kafka 吃磁盘 IOClickHouse 吃内存混在一起会互相影响。部署的第一步是准备 JDK。Flink 1.13 以上的版本要求 JDK 8 或者 11我建议用 JDK 11避免后面跑高版本连接器时出现模块访问限制。下载完 Flink 后需要修改 conf/flink-conf.yaml 里的三个关键参数jobmanager.memory.process.size 设置 2gtaskmanager.memory.process.size 设置 8gtaskmanager.numberOfTaskSlots 设置 4。这里有一个容易搞错的地方taskmanager.memory.process.size 不是堆内存而是整个进程的内存预算包括堆外内存和 JVM 元空间所以不要试图把它全部配成 -Xmx 的大小否则容器 OOM 会来得非常快。jobmanager.memory.process.size: 2048m taskmanager.memory.process.size: 8192m taskmanager.numberOfTaskSlots: 4 parallelism.default: 2这段配置的逻辑是一个 TaskManager 进程内部开多个 slot每个 slot 跑一个任务线程。parallelism.default 设为 2表示默认并行度为 2如果数据源有 4 个分区而并行度只有 2那有两个分区的数据会排队吞吐上不去。我通常建议把并行度和 Kafka 分区数设为一致或者并行度是分区数的整数倍这样最省心。部署完 Flink 后启动脚本是 bin/start-cluster.sh然后用 jps 看进程确认 StandaloneSessionClusterEntrypoint 和 TaskManagerRunner 都在跑再访问 8081 端口看 Web UI。3.2 ClickHouse 21.8 安装包与配置关闭 swap、调大 max_memory_usageClickHouse 的安装在 Linux 上非常快。21.8 LTS 版本的 rpm 包安装后默认数据目录在 /var/lib/clickhouse配置文件在 /etc/clickhouse-server。安装完成后第一件事不是启动而是改两个地方关闭操作系统的 swap因为 ClickHouse 在内存不足时会疯狂使用 swap 导致查询延迟飙升修改 users.xml 里的 max_memory_usage单机测试时我一般调到物理内存的 60%比如 32G 机器就设 20G这个参数决定了单个查询最多能用多少内存。profiles default max_memory_usage20000000000/max_memory_usage max_memory_usage_for_all_queries30000000000/max_memory_usage_for_all_queries max_partitions_per_insert_block1000/max_partitions_per_insert_block /default /profiles这里的逻辑是max_memory_usage 限制单个查询的内存max_memory_usage_for_all_queries 限制整个实例上所有并发查询的内存总和。电商看板场景里经常出现多个报表同时刷新的情况如果只配了前者两个大查询同时跑还是会把内存打爆。max_partitions_per_insert_block 是防止一次插入的 block 里包含太多分区导致 merge 任务积压默认 100如果你的写入任务按天分区但一次写入跨了三十天这个值够用如果 Flink 任务里有历史数据回填跨几百个分区这里就得调大。3.3 Kafka 在这里不是可选项它是 Flink 的削峰缓冲层有些同学会问MySQL 的 binlog 直接通过 Canal 推到 Flink 不行吗为什么要多一层 Kafka原因是 Flink 的 checkpoint 机制需要数据源支持回放而 Kafka 天然支持按 offset 回放。当 Flink 任务失败重启时它能从最近一次 checkpoint 记录的 offset 继续消费中间没有消费到的数据不会丢。如果直接把 Flink 接到 CanalCanal 本身没有长期持久化能力一旦 Flink 任务挂掉重新拉起的瞬间Canal 里的数据可能已经过期。Kafka 在这条链路里的角色就是给数据流一个缓冲区和后悔药。Kafka 部署时我建议用 KRaft 模式而不是 Zookeeper 模式21.8 版本的生态已经完全兼容 KRaft省掉一套 ZK 进程大幅降低运维负担。主题创建时实时订单流水建议分 6 到 12 个分区分区太少会导致 Flink 并行度上不去分区太多又会让单条消息的延迟增加因为每个分区在 Kafka 内部是一个文件目录太多小文件会让磁盘随机读写变多。订单流水这类消息我一般设置 7 天保留期点击流埋点数据量更大设置 3 天就够了因为实时分析只关心最近一段时间的窗口数据历史数据已经由离线数仓接管了。4. 用 Flink 实现 MySQL 同步到 ClickHouse从 CDC 捕获到 JDBC 写入的完整链路4.1 数据从哪来MySQL binlog 的 CDC 捕获与消息格式约定在这个项目里订单核心数据存储在 MySQL需要实时同步到 ClickHouse。常见做法是使用 Canal 或者 Debezium 监听 MySQL binlog把 insert、update、delete 操作解析成一条条变更记录写入 Kafka。我选择 Canal 比较多因为它在国内电商团队里普及度高部署文档齐全而且对 MySQL 主从同步协议的支持非常稳定。Canal 的配置里有一个关键点binlog 格式必须设置为 ROW因为 STATEMENT 格式记录的是 SQL 语句而不是数据变更前后的值无法支撑下游还原数据。Canal 投递到 Kafka 的消息格式通常是 JSON里面包含 data 字段、old 字段、type 字段、table 字段。data 是变更后的完整行数据old 是更新前的旧值type 区分 insert、update、delete。Flink 侧接入的时候需要特别注意 type 字段的判断。常见做法是解析 JSON 后维护一个表结构映射如果是 delete 事件不能直接把整行数据写入 ClickHouse因为 ClickHouse 的 MergeTree 引擎默认不支持删除操作需要把 delete 转换成 ReplacingMergeTree 里的一个状态标记或者使用 CollapsingMergeTree 把负向数据写入这样后续做 sum 聚合时正负抵消达到逻辑删除的目的。4.2 Flink JDBC 连接器配置批量写入与背压调优的四个参数Flink 官方 JDBC 连接器支持写入 ClickHouse但直接用它写默认是单条 insert吞吐低到没法看。实际落地时有两种做法一种是使用 flink-connector-clickhouse 第三方连接器它封装了批量插入和异步写入另一种是自己在 JDBC sink 上做 buffer 控制。我更倾向于先理解 JDBC sink 的批量机制再用第三方连接器这样出问题时不至于黑匣子。下面是用 DataStream API 构建 ClickHouse sink 的关键代码底层用的是 flink-connector-jdbc但把写入方式改成了手动批量 flush。public class ClickHouseSinkFunction extends RichSinkFunctionString { private static final String INSERT_SQL INSERT INTO order_flow (order_id, user_id, sku_id, pay_amount, event_time) VALUES (?, ?, ?, ?, ?); private Connection conn; private PreparedStatement ps; private ListString buffer; private int batchSize 1000; private long lastFlushTime System.currentTimeMillis(); private long flushIntervalMs 5000; Override public void invoke(String value, Context context) throws Exception { buffer.add(value); if (buffer.size() batchSize || System.currentTimeMillis() - lastFlushTime flushIntervalMs) { flush(); } } private void flush() throws SQLException { for (String row : buffer) { // 解析 JSON为每行数据 setString / setLong ps.addBatch(); } ps.executeBatch(); conn.commit(); buffer.clear(); lastFlushTime System.currentTimeMillis(); } }这段代码的逻辑是每条数据先进入 buffer不立刻写入数据库当 buffer 达到 1000 条或者距离上一次写入过了 5 秒触发一次批量写入。这里有四个参数直接影响性能batchSize 控制单次写入多少条太大会导致 ClickHouse 单次插入的数据块过大merge 线程压力高太小又失去批量意义我一般设置在 1000 到 5000 之间。flushIntervalMs 控制最长等待时间这是为了防止流量低谷时数据长时间积在内存里电商凌晨流量小5 秒不写入的话 ClickHouse 侧延迟会变大。preparedStatement 的 rewriteBatchedStatements 要设成 true否则 JDBC 驱动不会把多条 insert 合并成一条多 VALUES 语句batch 只是减少了交互次数没有真正减少 SQL 解析开销。4.3 ClickHouse 建表与 Flink 写入的 Schema 对齐一个字段类型不一致就全链路翻车Flink 写入 ClickHouse 最容易翻车的地方是字段类型对齐。MySQL 里的 decimal(10,2) 同步到 ClickHouse 如果映射成 Float64精度会丢失金额数据错一分钱财务那边直接炸锅。正确的映射是 decimal 对应 ClickHouse 的 Decimal(18, 2)int 对应 Int32 或 Int64varchar 对应 Stringdatetime 对应 DateTime 或 DateTime64。日期时间类型还有一个时区问题MySQL 的 datetime 不带时区ClickHouse 的 DateTime 默认按服务器本地时区存储如果 Flink 任务运行的容器是 UTC 时区写入后时间会差 8 小时。我通常统一约定事件时间字段全部用 DateTime64(3, Asia/Shanghai)在 Flink 侧用 withTimestampFormat 的时区参数指定为 UTC8。CREATE TABLE order_flow ( order_id UInt64, user_id UInt64, sku_id UInt64, pay_amount Decimal(18, 2), order_status String, event_time DateTime64(3, Asia/Shanghai), day_partition Date ) ENGINE ReplacingMergeTree(order_status) PARTITION BY day_partition ORDER BY (order_id, event_time)这张表使用 ReplacingMergeTree以 order_id 为业务主键version 列在这里用 event_time 代替。ClickHouse 的 ReplacingMergeTree 并不是在写入时去重而是后台 merge 时才根据 ORDER BY 字段保留最新一条数据所以 Flink 端写入时不能依赖这个引擎做严格去重否则查询未 merge 的分区时会看到重复记录。这算是 ClickHouse 新手最容易踩的一个坑看板里统计订单数突然变多不是算错是重复数据还没有被 merge 掉。要解决它查询时不直接 count 表而是加 FINAL 关键字或者使用 argMax 聚合函数。5. 部署和运行中的常见问题从 flink 的 jdbc 连接器异常到 ClickHouse 内存爆掉5.1 JDBC 连接器报 Connection is not available request timed out现象Flink 任务运行两三天后突然出现大量 JDBC 连接超时异常任务进入重启循环看 Web UI 显示 source 端和 sink 端背压都高得离谱。原因ClickHouse 默认的 max_connections 是 1024Flink 的 TaskManager 默认连接池里会不断创建新连接如果某个时刻写入并发高连接池没有及时释放ClickHouse 侧连接数会打满。但更隐蔽的原因在 JDBC 驱动clickhouse-jdbc 老版本默认每个 Connection 内部还有一层 http 连接池两层连接叠加导致连接数超预期。解决统一在 Flink 侧限制连接池大小设置连接最大空闲时间为 60 秒同时把 ClickHouse 的 max_connections 调大并保证 Flink 连接器里设置 socketTimeout 为 30000 毫秒。还有一个排查技巧异常信息如果带有 Too many simultaneous queries 字样说明不是连接数问题而是查询并发超过了 max_concurrent_queries需要调这个参数而不是连接池。5.2 ClickHouse 查询突然变慢磁盘 IO 100%罪魁祸首是分区过多现象看板页面从秒级响应变成十几秒点击查询时 ClickHouse 的 CPU 不高但 iowait 居高不下systemctl 看 ClickHouse 日志里全是 MergeSortingTransform 和 marks loading 相关的慢查询。原因Flink 写入任务的分区字段设计不合理。比如按小时分区写入了一个月的数据就会产生 720 个分区每次查询都要在各个分区目录里找数据合并线程也来不及处理。另一个原因是 ReplacingMergeTree 执行 merge 时如果单分区内数据块过多它需要读大量数据做排序磁盘 IO 瞬间拉满。解决分区尽量按天设置不要按更细粒度。如果业务上必须按小时查询可以每天按小时分区但配合 TTL 把七天前的旧分区提前合并成天级分区。同时检查 ClickHouse 的 merge 配置把 background_pool_size 从默认 16 调大一点merge 任务并发提高能显著降低合并延迟。这个坑在部署文档里一般不会写但大促前压测时一定会暴露。5.3 数据重复和数据丢失同时出现窗口统计结果对不上账现象实时大屏的支付金额比 MySQL 实际数据多出 2% 到 5%同时部分维度下数据又缺失。Flink 任务的 checkpoint 一直成功但结果就是不对。原因这是一个典型的至少一次与精确一次混淆的问题。Flink checkpoint 保证了算子状态的一致性但 JDBC sink 的写入如果发生在 checkpoint 之前任务重启时会从上一个 checkpoint 重新写入一数据造成重复。如果 ClickHouse 表用的是普通 MergeTree重复数据天然存在。丢失则是因为 Kafka 的 offset 提交与数据处理不在同一个事务里Flink 消费了数据但还没写入 ClickHouse 时任务挂了恢复后会跳过这些数据。解决首先保证 Flink 端 exactly-once 的实现这里不需要分布式事务而是把 Kafka offset 存储在 Flink 的 checkpoint 中并在 ClickHouse 表设计上做幂等。例如 ReplacingMergeTree 配合全局唯一的业务主键重复写入会被最终去重。对于丢失的场景要把 checkpoint 间隔调小从 60 秒调到 10 到 15 秒这样重启后重放的数据量小且不会造成大面积重复。5.4 时间字段差 8 小时大屏的零点峰值总是提前或延后现象实时大屏在 23:00 到 01:00 之间波动明显和真实的电商促销订单时间对不上看起来像数据延迟但实际上是时区问题。原因Flink 任务运行在 Docker 容器里默认时区是 UTC事件时间被当成 UTC 解析。ClickHouse 表的 DateTime 字段又只用字符串存储查询时按服务器的本地时区展示两个环节各差八小时数据自然就错了。解决Flink 侧在用 from_unixtime 或者 SimpleDateFormat 解析埋点时间戳时强制指定 Asia/Shanghai 时区不要依赖操作系统默认时区。ClickHouse 建表时对时间字段使用 DateTime64(3, Asia/Shanghai)并在连接参数里加 use_server_time_zonefalse 和 server_time_zoneAsia/Shanghai。检查时用一条 SQL 验证select toTimeZone(event_time, Asia/Shanghai) from order_flow limit 1如果结果和 MySQL 里的原始时间一致说明时区链路没有出错。6. 进阶玩法把 Flink SQL 的维表 join 和 ClickHouse 的物化视图真正用起来部署和排错搞定之后这个项目的价值才开始释放。我建议下一阶段做三个进阶动作把平台从“能跑”变成“好用”。第一个是维表 join。电商实时分析里经常要在订单流里补充商品名称、用户城市等维度信息不可能每次查 ClickHouse 时再去关联维度表。把商品维表加载到 Flink 的 JVM 缓存里用 broadcast 流广播到所有任务节点订单流每来一条就可以直接本地命中joins 性能提升非常明显。维表变化不频繁的场景用定时刷新缓存即可注意别在缓存里放全量用户表亿级用户会直接把 TaskManager 堆内存撑爆。第二个是 ClickHouse 物化视图。Flink 已经算过的分钟级聚合结果写入一张明细表后ClickHouse 仍然可以用物化视图继续做小时级、天级累加。物化视图在 ClickHouse 里是插入时触发的数据写入源表时自动更新目标聚合表不需要 Flink 再起一个任务去消费 Kafka减少一条链路就少一个故障点。比如订单大屏需要小时级 GMV可以在明细表上建物化视图group by 小时、类目、渠道这样 Flink 只负责分钟级明细小时级结果完全由 ClickHouse 自身承担。第三个是查询侧加速。亿级数据直接查询时用 ORDER BY 对查询字段做优化设计把 where 里最常见的过滤字段放在 ORDER BY 靠前的位置。比如按天分区后经常按店铺过滤就把 shop_id 放在 ORDER BY 的第一个字段ClickHouse 查询时能快速跳过大量数据块。再配合跳数索引对高基数的 sku_id 建 bloom_filter 索引对低基数的 order_status 建 set 索引性能能再翻一两倍。这套平台做下来我最大的教训是不要把实时链路当成黑匣子任何一层出问题都不会报错给你看只会让最终数据变得不对。数据量小的时候感觉不到等亿级数据跑起来每个配置项都是在给未来的自己减负。希望帮到你。本文还有配套的精品资源点击获取
返回列表