简介:这是一套基于Flink与ClickHouse构建的亿级电商实时数据分析平台源码,覆盖PC端、移动端与小程序三端场景,适合计算机相关专业学生、教师及企业开发人员用于毕业设计、课程设计或项目立项演示。资源包共1136个文件,约7.07MB,以Java后端代码、JavaScript与Vue前端脚本、CSS样式、HTML页面及PNG图片资源为主,另含Markdown说明文档、XML配置、JSON数据与properties配置文件,前后端结构完整。项目已通过导师评审,答辩成绩达95分,代码经测试可正常运行。读者可获得完整的三端源码、部署文档与配套资料,便于理解实时数据采集、计算与可视化链路,也可在此基础上二次开发扩展功能。目前已有94人学习关注,适合作为大数据方向进阶实践与毕设参考。
1. 从一份「高分项目」说起:Flink+ClickHouse 到底在电商里扛什么活
电商后台的实时大屏,最怕的不是数据少,而是数据来得又快又乱。用户在 PC 端加购、在移动 App 下单、在小程序里领券,三条链路的埋点格式各不相同,峰值时每秒几万条事件涌进来,运营还要求「下单后 3 秒内看到分渠道 GMV」。这套基于 Flink+ClickHouse 的亿级电商实时数据分析平台,解决的就是这件事:Flink 负责把多端埋点清洗、聚合、对齐时间窗口,ClickHouse 负责把结果压成列存、支撑亚秒级的多维查询。它适合谁?适合手里已经有埋点日志、但还在用离线 T+1 出报表的团队,也适合想把这套架构当成课程设计或高分项目落地的同学。源码、部署文档、全套资料的价值不在于「能跑」,而在于你能顺着它看清一条事件从产生到可查的完整链路,以及每个环节该调哪个参数。
2. 架构拆解:Flink 到 ClickHouse 的数据流怎么设计才不堵
2.1 三层链路:采集层、计算层、存储层的职责边界
先把整条链路拆成三段看,后面调参才不会乱。采集层通常是埋点 SDK 把 PC、移动、小程序的事件统一成 JSON,落到 Kafka 的多个 topic,比如ods_page_view、ods_order、ods_cart。计算层是 Flink 作业,消费 Kafka 后做维度补全(把 user_id 关联出渠道、地区)、窗口聚合(1 分钟滚动窗口算 GMV)、去重(同一订单多次回调只算一次)。存储层是 ClickHouse,宽表按天分区,物化视图预聚合常用维度。
这里最容易翻车的是职责越界:有人把去重逻辑塞进 ClickHouse 的ReplacingMergeTree,结果查询时数据还没合并,指标忽高忽低。正确做法是 Flink 侧用keyBy(order_id)加状态去重,ClickHouse 只做存储和查询加速。选型理由也在这:Flink 有状态、能容错,ClickHouse 列存、写入快但更新弱,两者互补而不是互相替代。
2.2 用 Flink SQL 写一个最小可跑的聚合作业
如果不想一上来就写 DataStream,Flink SQL 是最快验证链路的方式。下面这段作业从 Kafka 读订单事件,按渠道和 1 分钟窗口聚合 GMV,再写入 ClickHouse。
-- 源表:Kafka 中的订单事件,JSON 格式 CREATE TABLE ods_order ( order_id STRING, user_id STRING, channel STRING, -- pc / app / miniapp amount DECIMAL(10,2), event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'ods_order', 'properties.bootstrap.servers' = 'kafka:9092', 'properties.group.id' = 'gmv_job', 'format' = 'json', 'scan.startup.mode' = 'group-offsets' ); -- 结果表:ClickHouse 宽表 CREATE TABLE ads_gmv_1min ( channel STRING, window_start TIMESTAMP(3), window_end TIMESTAMP(3), gmv DECIMAL(18,2), order_cnt BIGINT ) WITH ( 'connector' = 'clickhouse', 'url' = 'jdbc:clickhouse://clickhouse:8123/ads', 'table-name' = 'ads_gmv_1min', 'username' = 'default', 'password' = '', 'sink.batch-size' = '1000', 'sink.flush-interval' = '3s' ); -- 聚合逻辑:1 分钟滚动窗口,按渠道分组 INSERT INTO ads_gmv_1min SELECT channel, TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start, TUMBLE_END(event_time, INTERVAL '1' MINUTE) AS window_end, SUM(amount) AS gmv, COUNT(order_id) AS order_cnt FROM ods_order GROUP BY channel, TUMBLE(event_time, INTERVAL '1' MINUTE);逻辑说明:WATERMARK设 5 秒,是为了容忍移动端网络抖动导致的乱序,设太小会丢迟到数据,设太大窗口结果延迟高。scan.startup.mode用group-offsets而不是earliest,避免重启后重复消费历史数据把 GMV 算爆。参数上,sink.batch-size和sink.flush-interval是写入 ClickHouse 的关键:批量太小写入频繁,ClickHouse 会因小 part 过多触发合并压力;批量太大则延迟高。1000 条 / 3 秒是常见起步值,压测后再调。
2.3 ClickHouse 建表:分区键和排序键怎么定
ClickHouse 的查询性能八成取决于建表时的PARTITION BY和ORDER BY。电商实时分析最常见的查询是「按天 + 按渠道 + 按小时」,所以分区用天,排序键把渠道和时间放前面。
CREATE TABLE ads.ads_gmv_1min ( channel String, window_start DateTime, window_end DateTime, gmv Decimal(18,2), order_cnt UInt64, dt Date DEFAULT toDate(window_start) ) ENGINE = MergeTree PARTITION BY dt ORDER BY (channel, window_start) TTL dt + INTERVAL 90 DAY;逻辑说明:PARTITION BY dt让每天的数据独立成目录,查询带日期条件时能直接裁剪分区。ORDER BY (channel, window_start)让「某渠道某时间段」的查询走主键索引,避免全表扫描。TTL自动清理 90 天前数据,电商明细一般不需要永久保留,冷数据可以另存对象存储。注意Decimal精度要覆盖金额,UInt64用于计数,别用String存数字,否则聚合时类型转换会拖慢查询。
3. 部署落地:从零把平台跑起来的关键步骤
3.1 环境准备与组件版本对齐
部署文档里最容易忽略的是版本兼容。Flink 和 ClickHouse 的 JDBC 连接器版本必须匹配,否则会出现「连接器异常」这类玄学问题。常见组合是 Flink 1.17/1.18 配flink-connector-clickhouse对应版本,ClickHouse 用 23.x 以上。Kafka 用 3.x,ZooKeeper 可以省掉,直接用 KRaft 模式。
# 目录规划:所有组件放 /opt,数据放 /data mkdir -p /opt/{flink,clickhouse,kafka} /data/{clickhouse,kafka} # 启动 ClickHouse(单机版) cd /opt/clickhouse ./clickhouse server --config-file=config.xml & # 启动 Kafka(KRaft 模式,无需 ZooKeeper) cd /opt/kafka bin/kafka-storage.sh format -t $(bin/kafka-storage.sh random-uuid) -c config/kraft/server.properties bin/kafka-server-start.sh -daemon config/kraft/server.properties # 提交 Flink SQL 作业 cd /opt/flink bin/sql-client.sh -f /opt/jobs/gmv_job.sql逻辑说明:ClickHouse 单机版适合验证,生产要配副本和分片。Kafka KRaft 模式省去 ZooKeeper 运维,但要注意server.properties里的log.dirs指向/data/kafka,别放系统盘。Flink 提交 SQL 前,要把 ClickHouse 连接器 jar 放进lib/目录,否则会报ClassNotFoundException。
3.2 多端埋点数据对齐:PC、移动、小程序的字段映射
PC、移动、小程序三端的埋点字段名往往不一致,比如 PC 叫user_id,小程序叫openid,移动叫device_id。Flink 作业里要做一层字段映射,统一成user_id和channel。
-- 用 CASE WHEN 做渠道归一,用 COALESCE 做用户标识兜底 CREATE VIEW unified_order AS SELECT order_id, COALESCE(user_id, openid, device_id) AS user_id, CASE WHEN source = 'web' THEN 'pc' WHEN source = 'ios' THEN 'app' WHEN source = 'android' THEN 'app' WHEN source = 'wxapp' THEN 'miniapp' ELSE 'unknown' END AS channel, amount, event_time FROM ods_order_raw;逻辑说明:COALESCE按优先级取第一个非空值,保证用户标识尽量不丢。CASE WHEN把多端来源映射成统一渠道,后续聚合和查询都基于这个口径。注意unknown渠道要监控,如果占比突然升高,说明埋点字段变了或 SDK 版本不兼容,这是血泪经验——曾经因为小程序改字段名,导致整个渠道 GMV 归零,排查了半天才发现是映射没覆盖。
3.3 资源参数:Flink 并行度和 ClickHouse 写入线程怎么配
Flink 并行度不是越大越好。并行度受 Kafka 分区数限制,如果 Kafka topic 只有 6 个分区,Flink 并行度设 12 会有 6 个 slot 空转。常见做法是并行度等于 Kafka 分区数,或者略小于。ClickHouse 写入线程用sink.batch-size和sink.flush-interval控制,同时 ClickHouse 侧max_insert_threads别设太高,否则小 part 太多。
| 参数 | 建议值 | 说明 |
|---|---|---|
| Flink 并行度 | = Kafka 分区数 | 避免 slot 空转 |
| sink.batch-size | 1000~5000 | 太小写入频繁,太大延迟高 |
| sink.flush-interval | 3s~10s | 与 batch-size 配合 |
| ClickHouse max_insert_threads | 2~4 | 太高导致 part 过多 |
| ClickHouse max_memory_usage | 物理内存 70% | 留余量给系统 |
逻辑说明:这些值不是固定的,要根据实际压测调。比如峰值 QPS 高时,batch-size 可以到 5000,但 flush-interval 要相应缩短,否则数据在内存里积压。ClickHouse 的max_insert_threads设太高,写入并发大,但每个 insert 生成一个 part,后台合并跟不上就会报「too many parts」。
4. 避坑排查:Flink+ClickHouse 落地时最容易翻车的 5 个点
4.1 现象:作业重启后 GMV 翻倍
原因:Kafka 消费位点没提交成功,或者scan.startup.mode设成了earliest,重启后从头消费。解决:确认 checkpoint 开启,scan.startup.mode用group-offsets,并在 ClickHouse 侧用ReplacingMergeTree或 Flink 侧状态去重兜底。
4.2 现象:ClickHouse 报「Too many parts」
原因:写入批量太小或频率太高,每次 insert 生成一个 part,后台合并速度跟不上。解决:调大sink.batch-size,调长sink.flush-interval,同时降低max_insert_threads。如果已经积压,手动OPTIMIZE TABLE合并,但别频繁执行。
4.3 现象:Flink 作业反压,延迟越来越高
原因:ClickHouse 写入慢导致 sink 反压,或者聚合窗口状态太大。解决:先看 ClickHouse 写入耗时,如果慢就调批量参数;如果状态大,检查keyBy的 key 是否基数过高,比如用user_id做 key 会导致状态爆炸,应该用order_id或渠道。
4.4 现象:多端数据时间戳不一致,窗口结果对不上
原因:PC、移动、小程序的时间戳格式或时区不同,有的用毫秒,有的用秒,有的用 UTC。解决:在 Flink 里统一转成TIMESTAMP(3)并指定时区,WATERMARK要基于事件时间而不是处理时间。常见做法是在源表里用TO_TIMESTAMP_LTZ转换。
4.5 现象:ClickHouse 查询慢,大屏加载超时
原因:查询没走主键索引,或者分区裁剪失效。解决:确认查询条件带dt和channel,ORDER BY的前缀要匹配查询条件。如果还是慢,建物化视图预聚合,把常用维度的结果提前算好。注意物化视图会额外占存储,别滥用。
5. 进阶技巧:用物化视图和 Flink 状态 TTL 把平台压得更稳
5.1 物化视图预聚合:把大屏查询从秒级压到毫秒级
当数据量到亿级,即使 ClickHouse 列存快,实时大屏每次查全量聚合也会吃力。物化视图的思路是:在写入时就把常用维度的聚合结果算好,查询时直接读小表。
-- 按渠道和小时的预聚合物化视图 CREATE MATERIALIZED VIEW ads.mv_gmv_hourly ENGINE = SummingMergeTree PARTITION BY toDate(window_start) ORDER BY (channel, window_start) AS SELECT channel, toStartOfHour(window_start) AS window_start, sum(gmv) AS gmv, sum(order_cnt) AS order_cnt FROM ads.ads_gmv_1min GROUP BY channel, toStartOfHour(window_start);逻辑说明:SummingMergeTree会在后台自动合并相同 key 的行,把gmv和order_cnt累加。查询时直接SELECT * FROM mv_gmv_hourly WHERE channel='app',数据量比原始表小几个数量级。注意物化视图只对插入的数据生效,历史数据要手动INSERT INTO ... SELECT回填。
5.2 Flink 状态 TTL:防止去重状态无限膨胀
去重逻辑如果用keyBy(order_id)加状态,状态会随订单量无限增长,最终撑爆内存。必须设状态 TTL,让过期订单的状态自动清理。
-- Flink SQL 中设置状态 TTL(需在 TableConfig 中配置) -- 或在 DataStream API 中: StateTtlConfig ttlConfig = StateTtlConfig .newBuilder(Time.hours(24)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); valueStateDescriptor.enableTimeToLive(ttlConfig);逻辑说明:TTL 设 24 小时,是因为订单回调一般不会超过一天,超过的重复回调概率极低。NeverReturnExpired保证过期状态不会被读到,避免脏数据。注意 TTL 不是立即清理,而是惰性清理,内存紧张时可以配合cleanupInRocksDBCompactFilter加速。
5.3 一个验证方法:用对账脚本核对实时和离线结果
实时平台最怕「看起来对,其实错」。我一般会写一个对账脚本,每天凌晨把 ClickHouse 的实时结果和离线 Hive 的结果按渠道对比,差异超过 0.5% 就告警。
# 对账脚本:对比实时和离线 GMV import clickhouse_driver client = clickhouse_driver.Client(host='clickhouse') realtime = client.execute(""" SELECT channel, sum(gmv) FROM ads.ads_gmv_1min WHERE dt = today() - 1 GROUP BY channel """) offline = client.execute(""" SELECT channel, sum(gmv) FROM offline.dws_gmv_daily WHERE dt = today() - 1 GROUP BY channel """) for ch, rt_gmv in realtime: off_gmv = dict(offline).get(ch, 0) diff = abs(rt_gmv - off_gmv) / max(off_gmv, 1) if diff > 0.005: print(f"渠道 {ch} 差异 {diff:.2%},实时 {rt_gmv},离线 {off_gmv}")逻辑说明:对账是最后一道防线,能发现窗口丢数、去重失效、字段映射错误等问题。差异阈值 0.5% 是经验值,太严会频繁告警,太松会漏掉问题。脚本要每天跑,结果存档,方便回溯。
这套平台我踩过的最大坑,是早期没设状态 TTL,跑了三天后 Flink 作业 OOM,重启后状态恢复又花了半小时。后来养成习惯:任何带状态的作业,上线前先问自己「状态会不会无限增长」。希望帮到你。
本文还有配套的精品资源,点击获取