简介:这是一份基于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 吃磁盘 IO,ClickHouse 吃内存,混在一起会互相影响。
部署的第一步是准备 JDK。Flink 1.13 以上的版本要求 JDK 8 或者 11,我建议用 JDK 11,避免后面跑高版本连接器时出现模块访问限制。下载完 Flink 后,需要修改 conf/flink-conf.yaml 里的三个关键参数:jobmanager.memory.process.size 设置 2g,taskmanager.memory.process.size 设置 8g,taskmanager.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_usage
ClickHouse 的安装在 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_usage>20000000000</max_memory_usage> <max_memory_usage_for_all_queries>30000000000</max_memory_usage_for_all_queries> <max_partitions_per_insert_block>1000</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 接到 Canal,Canal 本身没有长期持久化能力,一旦 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 RichSinkFunction<String> { 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 List<String> 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 或 Int64,varchar 对应 String,datetime 对应 DateTime 或 DateTime64。日期时间类型还有一个时区问题:MySQL 的 datetime 不带时区,ClickHouse 的 DateTime 默认按服务器本地时区存储,如果 Flink 任务运行的容器是 UTC 时区,写入后时间会差 8 小时。我通常统一约定:事件时间字段全部用 DateTime64(3, 'Asia/Shanghai'),在 Flink 侧用 withTimestampFormat 的时区参数指定为 UTC+8。
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 是 1024,Flink 的 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_zone=false 和 server_time_zone=Asia/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 索引,性能能再翻一两倍。
这套平台做下来,我最大的教训是不要把实时链路当成黑匣子,任何一层出问题都不会报错给你看,只会让最终数据变得不对。数据量小的时候感觉不到,等亿级数据跑起来,每个配置项都是在给未来的自己减负。希望帮到你。
本文还有配套的精品资源,点击获取