接手营销自动化平台的数据架构时,我们面对的是一堆散落在各个业务库里的用户行为、订单、投放回执和客服记录,运营同学每次做活动复盘都要在 Excel 里手工拉数。后来我们把多源数据汇总进统一的 OLAP 分析引擎,营销活动的定向圈人、漏斗分析、效果归因才真正跑了起来。这篇文章就来复盘这套多源数据 OLAP 架构的演进过程,聊聊我在选型和落地时踩过的坑,以及每一版方案背后到底在解决什么问题。
营销自动化数据驱动:多源数据 OLAP 架构演进实录
1. 项目背景:营销自动化场景下的数据困境
营销自动化这个概念听着很洋气,落到实际业务里无非就是三件事:在合适的时机,通过合适的渠道,把合适的内容推给合适的人。听起来简单,但“合适”这两个字背后全是数据问题。
1.1 营销平台为什么离不开数据分析
拿我们当时的一个典型场景来说:运营同学要做一个“近30天加购但未下单,且最近7天活跃过的女性用户”的召回活动。这个人群定义涉及四个不同的判断条件——加购行为、下单状态、活跃时间、用户性别——而这四类数据分别储存在行为日志库、交易订单库和用户画像库里。
在没有统一分析层之前,运营同学取数靠的是“求人”。先找数仓组写 Hive 脚本跑离线任务,跑完导出 CSV,再导入 Excel 做透视表,最后才能上传到营销后台做圈选。一套流程下来,一个简单的人群定义要等两三天,活动节点早就过了。
而且营销活动一旦跑起来,效果分析又是新一轮痛苦。短信发了50万条,哪个渠道带来的转化最高?不同人群包的 ROI 分别是多少?优惠券核销和用户留存之间是什么关系?这些问题全部需要将触达数据、行为数据、交易数据三者关联分析,单一的业务系统数据库根本扛不住这种分析负载。
1.2 数据团队面临的三个转型节点
我在接手数据架构的两年多时间里,整个系统经历了三次明显的转型节点。
第一个节点是从无到有。我们最早只有业务库(MySQL)直连 BI 报表,随着数据量增长,大表关联查询能把主库拖到报警。于是开始搭建离线数仓,把历史数据同步到 Hive,形成 ODS/DWD/DWS 分层。
第二个节点是从慢到快。离线数仓解决了主库压力,但 T+1 的时效性让营销活动根本没法“自动化”。今日的实时推荐、实时人群包更新、实时活动效果监控,都要求查询响应从小时级压缩到秒级。这时候OLAP 引擎(ClickHouse / Doris 这类列式存储分析数据库)就浮出水面了。
第三个节点是从单点到多点。数据源从最初的订单和埋点日志,扩展到广告投放回传、企微会话存档、客服工单、App 推送回执等十几个来源,实时与离线链路并行。多源数据的标准统一、标识映射、时序对齐成了新难题。
这篇文章记录的核心,就是我们在第三个节点前后如何把 OLAP 架构从 V1 演进到 V3 版本的整个过程。
2. 多源数据接入:先理顺源头才能算得准
OLAP 引擎只是“加工车间”,真正决定分析结果上限的是进来的“原材料”。营销场景的数据源我梳理了一下,至少有六大类,每类的接入方式和分析价值都不太一样。
2.1 营销数据源的六大分类
根据我实际接触到的业务,可以把营销相关的数据源归成下面这张表:
| 数据源类型 | 典型来源 | 核心数据 | 更新频率 | 分析价值 |
|---|---|---|---|---|
| 用户行为数据 | App/H5埋点、小程序日志 | 浏览、点击、加购、收藏 | 实时/准实时 | 兴趣识别、漏斗分析 |
| 交易订单数据 | 交易库、订单中心 | 订单金额、商品明细、支付状态 | 近实时/离线 | GMV分析、转化归因 |
| 触达投放数据 | 短信平台、Push服务、投放渠道 | 发送记录、送达状态、回执 | 实时/离线 | 渠道效果评估 |
| 客户管理数据 | CRM、私域运营系统 | 客户等级、归属导购、标签 | 离线/准实时 | 客户分群、价值分层 |
| 客服互动数据 | 客服工单、会话存档 | 咨询记录、满意度、投诉类型 | 离线 | 体验分析、流失预警 |
| 广告投放数据 | 媒体渠道API回传 | 曝光、点击、消耗、转化回传 | 实时 | ROI计算、投放优化 |
每个数据源都有自己的“脾气”。行为日志天然是流式的,按时间堆积且量级最大;订单数据有状态流转(待支付→已支付→已退款),同步时要特别注意状态覆盖;投放回传数据则是典型的“晚到数据”,经常出现延迟三五个小时才回来的转化,需要设计好回刷机制。
2.2 身份标识统一:多源数据关联的第一道坎
多源数据接入时最容易翻车的,就是不同系统里的同一个用户根本长得不一样。
登录用户在行为日志里带 user_id,在订单库里带 uid 关联到用户表,在 CRM 里存的是手机号,广告回传就更复杂——媒体回传的是 IDFA / OAID / 广告ID,和业务库里的用户标识完全不沾边。
我们当时的做法是建立一套统一身份域概念,核心是三层映射:
- 设备层映射:device_id / IDFA / OAID → 设备指纹
- 账号层映射:user_id / email / phone → 统一账号
- 业务层映射:统一账号 → 客户ID / 会员ID / 导购ID
落地上,我在 ODS 层加了一个用户标识映射表(user_mapping),不断吸收各数据源产生的关联关系,再用离线任务+实时 UPSERT 的方式维护映射关系。每个源表进来以后,统一替换成标准 user_key 字段再进入 DWD 层。
提示:如果不做标识映射就直接拼接多源数据,后期做跨渠道去重和频控时,你会被重复计算折磨到崩溃。这一步一定要前置,越早越好。
2.3 数据接入技术选型:离线与实时的双通道
数据接入我们最终采用了离线批量 + 实时流式双通道的方案,简单说就是“同样一双鞋,白天走楼梯,晚上坐电梯”:
- 离线通道:基于 Apache Airflow 调度,通过 DataX / SeaTunnel 把各业务库数据每日全量或增量同步到数据湖(HDFS/Iceberg)和数仓表中。负责补齐历史数据和修正异常数据。
- 实时通道:MySQL 的变更数据通过 Flink CDC 写入 Kafka,实时行为日志直接通过 SDK/日志采集器进 Kafka,再由 Flink 做 ETL 后写入 OLAP 引擎的实时表。
双通道并行带来一个必然问题:数据重复和冲突。我们的解决策略是给每条数据打上数据源优先级和事件时间戳(event_time),在同一主键冲突时,以优先级高、时间最新的数据为准。
营销场景实时链路示意:
# 各业务源(MySQL/埋点SDK) # ↓ # Flink CDC / 日志采集 # ↓ # Kafka(多topic,按业务域隔离) # ↓ # Flink ETL(清洗、标准化、user_key映射) # ↓ # ClickHouse / Doris 实时表很多团队一开始图省事,所有数据都走一条实时链路,结果 Kafka 堆积、OLAP 写入过载,反而比离线还慢。我的建议是:分清数据的“时效敏感性”,用户行为和投放回执必须实时,订单和 CRM 数据准实时(分钟级)就够,历史数据用离线补。
3. OLAP 架构演进路线:从单机 MySQL 到湖仓一体分析
很多文章讲架构演进喜欢直接给终态方案,我觉得这反而对读者没价值。这里我把我们实际的演进过程分成三个版本,说说每个版本遇到什么问题、为什么撑不住了,才是真正的“演进”。
3.1 V1 阶段:单库直连 + 定时报表的原始时代
最开始我们根本没有独立的数据平台。运营要看数据,直接连 MySQL 业务库写 SQL,或者 dumps 数据到本地跑 Python 脚本。
痛点非常明显:
- 性能瓶颈:一张订单明细表 5000 万行,关联用户表做聚合查询,一个 count + group by 的简单统计能把主库 CPU 打到 80% 以上。营销活动一做,正常交易都受影响。
- 无统一口径:不同运营从不同表口径取数,加购率有的算成加购人数/访客数,有的算成加购件数/浏览量,一份周报对不上,市场部和运营部先在群里打一架。
- 时效性差:报表靠人工触发,领导早上要的数据,下午才跑完。
V1 阶段唯一的收获,是帮我们验证了一件事——营销数据分析的需求量远超预期。每个活动都是一轮从圈选到复盘的全流程取数,没有一套专门的分析引擎是根本扛不住的。
3.2 V2 阶段:离线数仓 + Presto/Spark 的标准化时代
V1 撑到业务规模上来之后,我们开始搭建规范的离线数仓。以 Hive 为核心,按 ODS(原始数据层)、DWD(明细数据层)、DWS(汇总数据层)、ADS(应用数据层)四层模型组织数据。
- ODS 层:直接同步各源系统的原始数据,保留全量历史。
- DWD 层:清洗、去重、维度退化、统一用户标识,生成标准明细事实表。
- DWS 层:按业务主题(用户、商品、活动)做轻度汇总,例如用户购买频次表、活动参与明细表。
- ADS 层:面向应用的多维汇总结果,直接服务 BI 报表和人群圈选接口。
这个阶段解决了“有没有标准和能不能跑”的问题,但新问题很快出现:
离线数仓的典型时效是 T+1,这意味着今天看昨天的数据。但对营销自动化来说,“昨天”太慢了。比如用户在 App 里把商品加入购物车,如果到第二天才被圈进“加购未购人群”,这个用户可能早就去别的平台下单了。
我们用 Presto 在 Hive 之上做快速查询,但也只是从“小时级”提升到“分钟级”,离营销场景的需求还是差一个量级。运营同学说得对:“我们要的不是昨天发生了什么,而是现在这一刻谁该触达。”
3.3 V3 阶段:实时 OLAP + 湖仓一体的自动化时代
为了把查询耗时控制在秒级,我们在应用层引入了实时 OLAP 引擎。架构上形成了一套“实时优先、离线兜底”的双轨体系:
- 实时链路:Kafka → Flink → ClickHouse 实时表,支撑实时人群圈选、实时活动监控、实时大屏。
- 离线链路:业务库 → 数仓(Hive/Iceberg)→ 批量写入 ClickHouse/BI,支撑深度分析、历史回溯、T+1报表。
- 统一分析入口:所有面向分析的服务都走统一的 OLAP 查询层,不再直连业务库。
为什么把 ClickHouse 放到应用层而不是替代数仓?因为两者解决的问题不同。Hive/Iceberg 存储成本低、事务能力强、适合大规模全量扫描计算;ClickHouse 列式存储、向量化执行、压缩比高,非常适合高并发多维聚合查询,但在事务写入和复杂更新上比较弱。用一套引擎通吃所有场景,本身就是一种不切实际的幻想。
V3 阶段之后,我们的整个技术栈形成了这样一张全景:
| 层次 | 组件 | 职责 |
|---|---|---|
| 数据源层 | 埋点SDK、订单库、CRM、广告回传、客服系统 | 产生原始业务数据 |
| 接入层 | Kafka、DataX/SeaTunnel、Flink CDC | 实时/离线采集传输 |
| 存储计算层 | Hive/Iceberg + Spark/Presto | 离线数仓、大规模批处理 |
| 实时处理层 | Flink + Kafka | 流式ETL、实时聚合 |
| 分析服务层 | ClickHouse(主)、MySQL(元数据) | 秒级OLAP、高并发查询 |
| 应用层 | 人群圈选平台、报表系统、营销自动化编排引擎 | 面向业务输出能力 |
这套架构跑下来,营销活动的圈选响应时间从平均 3 小时降到 300 毫秒,周报自动化率达到 95% 以上。
4. 关键查询场景与 OLAP 技术选型
架构演进不是拍脑袋决定的,每一步都是被具体的分析场景逼出来的。这一节我来拆解营销自动化平台里最常见的四类 OLAP 查询,以及它们在技术选型上对应的要求。
4.1 人群圈选:广告系统之外的“实时分群”难题
人群圈选是营销自动化的核心入口。运营同学在前端界面配置条件(如“最近30天加购2次以上且消费总额>500元的上海女性用户”),后端的查询引擎需要在秒级返回符合条件的人群包 ID 列表。
这类查询的核心特征是多条件组合过滤 + 精确去重 + 结果集返回。传统 SQL 写法是:
SELECT user_key FROM dws_user_behavior_agg WHERE 1=1 AND last_30d_add_cart_cnt >= 2 AND total_consumption_amount > 500 AND region = '上海' AND gender = 'female' AND last_active_day >= today() - 7如果 dws_user_behavior_agg 是一张大宽表(每个用户一行,包含各种指标字段),ClickHouse 这种列式存储就很适合——它只需要读取涉及的几列数据做过滤,存储层能跳过无关数据块。
另一个难点是人群包交集/并集/差集运算。一个活动可能要排除“近7天已触达过的用户”,或用 A 人群和 B 人群做交集。我用 ClickHouse 的 Bitmap 数据类型存储人群包,通过位图运算完成集合操作,性能非常好:
-- 从人群包表获取两个人群的 bitmap 并求交集 SELECT bitmapAndOrCardinality( (SELECT bitmap FROM audience_pack WHERE pack_id = 'A'), (SELECT bitmap FROM audience_pack WHERE pack_id = 'B') ) AS intersect_cnt;在千万级用户规模下,这种位图运算的耗时基本在毫秒级,比 JOIN 方式快几个数量级。
4.2 漏斗分析:把“用户走到哪一步流失”看明白
运营同学最喜欢看的是转化漏斗。比如“从启动 App 到完成支付”的漏斗:启动 → 浏览商品 → 加购 → 提交订单 → 支付成功。每一层的转化率像体温计一样反映活动效果。
漏斗分析在实现上有一个天然的麻烦——每一步事件可能分散在不同时间点,必须按用户做会话内的路径串联。在 ClickHouse 里,我用 windowFunnel 函数来做:
SELECT level, count() AS user_cnt FROM ( SELECT user_key, windowFunnel(3600)(event_time, event_type = 'app_launch', event_type = 'product_view', event_type = 'add_cart', event_type = 'submit_order', event_type = 'pay_success') AS level FROM dwd_user_event_log WHERE event_date = today() GROUP BY user_key ) GROUP BY level ORDER BY level;windowFunnel 函数接受一个时间窗口(秒)和一组事件条件,按顺序判断每个用户在窗口内是否依次完成了这些行为,返回命中的最大层级。这个函数处理几千万行明细数据的耗时也是秒级。
漏斗分析的另一个注意点是会话切分。用户可能上午浏览没买,晚上又回来直接下单了,中间隔了好几个小时,这算一个漏斗还是两个?我们当时定的规则是:相邻行为间隔超过 30 分钟,视为两个会话。这个参数可以做成配置项,不同业务线看自己的习惯调整。
4.3 留存分析与触达效果评估:衡量自动化真实价值
自动化营销上线之后,老板最关心的问题并不是“触达了多少人”,而是“触达之后用户留下了吗”。留存分析需要把用户按首次行为时间(通常是注册或首购)分桶,再统计后续每天的活跃/购买比例。
留存分析的 SQL 看起来简单,但数据量大的时候很容易卡死:
SELECT first_day, datediff('day', first_day, active_day) AS day_diff, uniqExact(user_key) AS retained_users FROM ( SELECT user_key, toDate(min(event_time)) AS first_day, toDate(event_time) AS active_day FROM dwd_user_event_log GROUP BY user_key, toDate(event_time) ) GROUP BY first_day, day_diff;在两三千万行明细上跑这种查询,ClickHouse 通常一秒内能出结果。但有几个细节值得注意:
一是去重口径。同一个用户同一天多条行为只算一个活跃用户,所以要用 uniqExact 而不是 count;二是日期维度的补齐。留存矩阵中最右边的格子可能没有数据,需要在应用层做笛卡尔积补齐前端的展示;三是留存口径定义要固化,我们把留存指标直接做成了 ClickHouse 的物化视图,应用层查视图而不是直接跑明细。
触达效果评估则涉及归因问题。用户收到短信后没有立刻下单,而是第二天主动打开 App 买了东西,这笔转化算不算短信的功劳?我们采取的是“首次触达归因 + 功劳衰减”的策略:用户 7 天内发生转化,按触达渠道的先后顺序分配权重,第一个触达渠道拿到最高的贡献比例。
归因计算放在 Flink 实时任务里做,每来一条转化事件就去查询该用户近7天的触达记录,计算完写入结果表。OLAP 引擎在其中起的作用是提供快速的触达历史查询接口。
4.4 ClickHouse 与 Doris 的选型对比
聊一下引擎选型。当时我们在 ClickHouse 和 Apache Doris 之间纠结了很久,做了三轮 POC 对比,最终选择了 ClickHouse。选型过程比结论更有参考价值,这里把我整理过的关键对比放出来:
| 维度 | ClickHouse | Apache Doris |
|---|---|---|
| 列式存储/向量化执行 | 成熟、性能极致 | 成熟,兼容 MySQL 协议 |
| 高并发查询 | 单机并发 100~200 优化后可达 1000+ | 并发能力较好,适合高并发点查 |
| JOIN 能力 | 弱于 Doris,需要大表驱动小表或提前宽表化 | 支持多种 JOIN,优化较好 |
| 精确去重 | Bitmap 类函数强大 | Bitmap 也有实现 |
| 运维复杂度 | 集群依赖 ZK,组件较重 | 相对轻量,FE/BE 架构清晰 |
| 实时写入 | Kafka 集成,批量写入吞吐高 | 支持实时导入,StreamLoad 很方便 |
| 生态成熟度 | 社区活跃,但函数/语法比较特立独行 | 更贴近常规 SQL,学习成本低 |
我们选 ClickHouse 的原因很务实:
- 我们的核心场景是大宽表 + 多维聚合 + 大结果集,这正是 ClickHouse 的长项;
- 团队里有人有 ClickHouse 生产经验,排障成本低;
- 人群圈选场景的 Bitmap 能力以及 windowFunnel 漏斗函数,在当时的 Doris 版本上不如 ClickHouse 成熟。
但如果你后期有很强的并发点查需求(比如 To C 场景的高并发在线服务),或者团队对传统 SQL 更熟悉,Doris 也是完全值得考虑的。没有脱离场景的最好引擎,只有最适合自己团队的选型。
5. 核心实现与调优细节:把 OLAP 榨出应有的性能
选型定了,真正的挑战才开始。接下来聊聊我们在 ClickHouse 上做表设计、实时链路和查询优化时的具体实现细节,以及为什么这么设计。
5.1 表引擎与分区设计:物理模型决定查询天花板
ClickHouse 的表引擎选择直接决定了查询性能的上限。我们的核心事实表(用户行为事件表、订单明细表)全部使用MergeTree 家族引擎:
- ReplacingMergeTree:用于需要去重的场景,例如 PAY 成功后重复回流的订单事件,用主键 + version 字段控制去重逻辑;
- SummingMergeTree:用于汇总表,相同的分组键会被自动累加,例如按小时聚合投放消耗的 DWS 表;
- 普通 MergeTree:用于明细数据,保留全量,查询时靠分区裁剪和索引加速。
分区设计这块踩过一个大坑。最早我们把事件表按toYYYYMM(event_time)(按月)分区,结果查询近三天的实时分析要扫描整整一个月的分区数据。后来改成按天分区:
ENGINE = MergeTree PARTITION BY toYYYYMMDD(event_time) ORDER BY (user_key, event_time)查询过滤条件带上具体日期时,分区裁剪可以让扫描量下降一个数量级。分区粒度粗了浪费 IO,细了(比如按小时)会产生大量小文件,合并压力大。根据我们的经验,营销行为日志按天分区是性价比最高的折中方案。
5.2 Bitmap 精确去重与正交分桶
营销分析里“去重”几乎无处不在:留存分析要算去重用户、投放去重要算去重设备、人群圈选要算去重身份。常见做法是uniqExact函数,但它在亿级数据量下会消耗大量内存。
我们在处理“亿级用户 + 多维组合”的计数场景时,采用Bitmap 预计算方案。思路是:提前把用户 ID 按照枚举维度(渠道、地区、注册月份等)构建成 Bitmap 存储:
CREATE TABLE audience_bitmap ( event_date Date, channel_id String, region_id String, user_bitmap Bitmap, ... ) ENGINE = MergeTree PARTITION BY event_date ORDER BY (channel_id, region_id);写入时用groupBitmapState(user_key)聚合函数构建位图,查询时直接用bitmapCount计算去重人数:
SELECT channel_id, bitmapCount(bitmapAnd(user_bitmap, (SELECT user_bitmap FROM audience_bitmap WHERE region_id = '上海'))) AS shanghai_cnt FROM audience_bitmap GROUP BY channel_id;这套方案在千万级用户上做任意维度的超高速去重统计,查询耗时基本在 100ms 以内。但注意它只适合“预聚合 + 枚举维度”的场景,如果是任意维度临时组合,还是要走明细聚合。
5.3 实时链路的数据延迟与幂等保障
实时链路最怕的是“算错、算重、算亏”。我们的实时架构是这样的:
Flink 消费 Kafka 数据 → ETL 清洗 → 写入 ClickHouse
在写入 ClickHouse 环节,我们用的不是 Flink JDBC Sink,而是Kafka 引擎表过渡 + 定时批量落盘的方案:
-- ClickHouse 中创建 Kafka 引擎表作为缓冲 CREATE TABLE kafka_user_event_buffer ( user_key String, event_type String, event_time DateTime, ... ) ENGINE = Kafka() SETTINGS kafka_broker_list = 'broker:9092', kafka_topic_list = 'dwd_user_event', kafka_group_name = 'clickhouse_consumer', kafka_format = 'JSONEachRow'; -- 物化视图将 Kafka 数据流转到实际存储表 CREATE MATERIALIZED VIEW mv_user_event TO user_event AS SELECT * FROM kafka_user_event_buffer;这种方式的好处是:Flink 只管把清洗后的数据扔到 Kafka,ClickHouse 自己消费并批量写入,减少了 Flink 到 ClickHouse 的直接连接压力,同时天然做成了解耦。
幂等保障是另一个容易被忽略的细节。Kafka 数据可能重复投递(查重算法用的是 at-least-once),所以我们在目标表上用了 ReplacingMergeTree 引擎,以user_key + event_time + event_type作为主键,重复数据会在后台合并时被清理。
这里提醒一句:ReplacingMergeTree 的去重不是实时的,合并发生在后台异步完成,所以刚写入的几分钟内可能查到重复数据。如果业务对实时去重有硬要求,可以在应用层用 bitmap 或 bloom filter 做二次检查。
5.4 宽表设计与维度退化:少 JOIN 才能快
ClickHouse 的 JOIN 能力不像传统数据库那么强,大表之间 JOIN 经常内存爆炸或跑成“全村吃饭慢查询”。我们的策略是尽量做宽表、减少 JOIN,也就是经典的维度退化思想。
举个例子:订单明细表早期只有order_id、user_id、shop_id、amount等 ID 字段,要分析“哪个城市的订单多”必须 JOIN 用户表拿城市。我们把用户的城市、性别、年龄、会员等级等常用维度字段直接冗余到订单明细表中:
-- 宽表化的订单事实表 CREATE TABLE dwd_order_wide ( order_id String, user_key String, order_amount Decimal(18,2), create_time DateTime, -- 冗余维度字段 user_city String, user_gender String, user_level String, pay_channel String, ... ) ENGINE = MergeTree PARTITION BY toYYYYMMDD(create_time) ORDER BY (user_key, create_time);宽表化的代价是存储空间增加和写入时 ETL 逻辑变重,但换来的是查询阶段几乎不需要 JOIN,性能提升是非常明显的。在营销自动化的大多数分析场景里,我强烈建议优先宽表、延后建模。
6. 演进过程中踩过的坑与排查实录
6.1 数据倾斜:热点用户把查询拖死
人群圈选上线后有一阵子经常出现“某个查询卡死 1 分钟以上”的情况,排查发现是数据倾斜问题。少部分活跃用户(比如头部主播的粉丝、大促期间疯狂下单的羊毛党)的行为数据量是普通用户的几百倍,导致按 user_key 分组聚合时某个分片处理的数据量远超其他分片。
几种解决思路:
- 加盐打散:在 user_key 后面拼接一个随机数把热点打散,计算完再汇总;
- 两阶段聚合:先按(user_key + 随机后缀)预聚合,再合并结果;
- 调整分片策略:ClickHouse 自定义分片键,把超大用户单独分桶;
- 业务层面限制范围:有的查询场景本身就不需要处理全量数据,加过滤条件把热点排除。
我们在实际落地时主要靠前两种:Flink 实时链路上的聚合任务在 keyBy 之前加了“热点检测 + 自动打散”的逻辑,离线聚合则在 SQL 层做两阶段优化。
6.2 时区不一致:埋点事件时间和业务时间差 8 小时
这是一个非常隐蔽但影响面极大的问题。我们埋点数据在客户端生成时用的是本地时区(东八区),而服务端订单时间存的是 UTC;刚开始直接把两边的 event_time 放进同一张表做分析,发现大部分“当天转化”的数据在时间维度上对不上。
排查思路不复杂,先查同一用户的行为事件和订单时间,拉了几条记录就发现了 8 小时差。后来我们统一了规范:所有进入 Kafka 的数据,在 Flink ETL 阶段统一转换为 UTC 存储,ClickHouse 表按 UTC 分区;查询展示层再做时区转换,BI 报表统一显示北京时间。这样避免了“同一时间戳在不同表里语义不一致”的数据人格分裂。
6.3 实时数据回刷:活动复盘和实时看板永远对不上
这是让运营和数据团队最抓狂的问题:实时看板显示 GMV 100 万,第二天离线报表跑完变 120 万。原因很简单——订单回流有延迟。支付成功事件不是实时的,微信/支付宝回调可能要几分钟甚至几小时才到达系统;凌晨大促时回调积压,更多数据在第二天才补录进来。
我们成立了一个“实时+离线一致性核对”机制:
- 每天凌晨离线任务把前一天的实时结果重新计算一遍,与实时看板做逐指标比对;
- 差异超过阈值(比如 3%)时,实时看板打上“数据修正中”的标记;
- 实时看板底部显示“数据更新至 xx:xx”,让管理层知道这是时间截点快照,不是最终数字。
顺带说一句:不要把实时链路设计成永远正确的方案。实时更像“快照”,离线才是“终版”。发布到高层决策的数据,永远以离线数据为准。
6.4 慢查询排查手册:从 SQL 到集群参数的定位路径
最后整理一份排查 OLAP 慢查询的路线图,遇到问题按这个顺序检查,大部分能快速定位:
| 排查步骤 | 检查内容 | 常用方法 |
|---|---|---|
| 1. 查询SQL本身 | 是否全表扫描?是否有 join 大表? | EXPLAIN 看执行计划 |
| 2. 分区裁剪 | where 条件是否带上了分区字段? | 看 read_rows / read_bytes |
| 3. 数据倾斜 | 某个分片 load 远高于其他分片? | 集群监控看每个 shard 的耗时 |
| 4. 索引是否生效 | 排序键与过滤条件是否匹配? | 检查 primary key 使用情况 |
| 5. 物化视图 | 高频查询是否可改写为物化视图? | 分析查询 pattern |
| 6. 资源配额 | 并发查询是否过多导致排队? | 查 system.processes / queue 时长 |
7. 给正在做类似架构演进的同学一些实在建议
如果你也在给营销平台搭建数据分析架构,我这里有一些基于实际踩坑提炼的建议,不写大道理,只写可以直接用的操作。
第一,不要一上来就上实时链路。先把离线数仓跑通,把数据口径、指标定义、质量校验机制都稳定下来,再考虑实时。实时链路是放大器,如果基础口径都是乱的,实时只会把乱放大得更快。
第二,建好“人、场、货”三个维度的基础宽表。营销分析 80% 的查询都是围绕这三个维度展开的,提前把行为、订单、触达数据组装成几张核心宽表,能避免后面大量重复开发。
第三,从第一天开始维护指标口径字典。什么算“活跃用户”?是启动过 App 就算,还是浏览过商品页才算?什么算“复购”?是再次下单还是再次支付?每个指标都写清楚定义、计算公式、来源表、同步频率。这件事不做好,数据产品根本无法推广到业务侧。
第四,把 OLAP 和业务系统做物理隔离。营销触达系统的实时决策查询和 BI 报表的历史分析查询负载特征完全不同,强行共用一套 OLAP 集群,必然互相影响。我们最终按业务优先级拆成了“实时决策集群”和“分析查询集群”,支持按资源配额隔离。
第五,重视数据治理,哪怕它看起来不重要。字段命名不统一、枚举值编码混乱、上游表结构变更没有通知——这些细碎的“脏乱差”会在 OLAP 层被成倍放大。我们在 Flink ETL 阶段写了一套 schema 校验 + 枚举值映射 + 异常数据旁路转储的机制,宁可丢弃异常数据也不要脏数据进入 OLAP 表。
在营销自动化这条路上,OLAP 架构演进没有标准答案,但方向是明确的:从离线到实时、从单源到多源、从被动查数到主动驱动业务决策。每一步演进都是为了回答同一个问题——如何更聪明地把对的内容送给对的人。希望这篇复盘能帮你少走几步弯路。