做营销自动化的同学应该都撞过同一堵墙:运营同事非常兴奋地跑过来说,"我想圈一下昨天加购但没付款、并且过去30天打开过App超过5次、最好还看过A商品详情页的用户,给他们推一张100减10的券。"听起来平平无奇对吧?但在你身后,事件明细表已经跑到了几十亿行,MySQL里那张带着大JSON字段的event表,一个GROUP BY就能让生产库CPU飘红。这不是段子,是我过去几年搭营销自动化数据体系时反复遇到的真实场景。
营销自动化的本质是一句话:在合适的时间,通过合适的渠道,把合适的内容推给合适的人。这句话里每一个"合适"背后,都藏着对多源数据的整合能力、对用户行为的实时理解能力、对历史效果的分析能力。这三种能力恰好都指向同一个底座——一套能扛住海量明细、支持多维分析、还能跟上实时节奏的OLAP架构。这篇文章从我亲身经历的角度,聊聊营销自动化场景下,多源数据OLAP架构从无到有、从简单到复杂的完整演进过程。
1. 营销自动化的OLAP需求:圈人、旅程、归因,哪一个都躲不开
1.1 运营同学的一句话需求,背后是海量明细的即时检索
前面那句"圈出昨天加购没付款、30天活跃5次、还看过A商品详情页的用户",翻译成技术语言大概是这么一回事:先从近30天的行为事件中,按user_id聚合统计每个用户的活跃天数;再过滤出昨天发生过add_to_cart事件、但后续没有payment事件的用户;然后在这些用户里找出近30天内浏览过特定商品详情页的那部分人;最后把满足三个条件的人群以列表形式返回给触达引擎。
这个查询在数据量小的时候写出来挺简单,三张表JOIN一下、几个WHERE条件一加就完事。但问题是,营销自动化业务跑起来之后,事件明细表是按天几千万行往上长的,累计几十亿行是常态。在这种量级下,MySQL行存引擎跑一次这样的聚合查询,磁盘IO直接被打满,一个查询跑十几分钟是家常便饭,更别提运营同学一次要圈好几批人、每批人数还不一样。
这就引出了OLAP在营销场景中的第一个核心价值:在海量明细上做即时检索和聚合。注意"即时"两个字——运营圈人不是在跑离线报表,他们拿到人群以后要立刻去配置触达任务,等不了半小时。所以查询性能必须控制在秒级甚至毫秒级。
1.2 营销自动化对OLAP的四个核心要求
我做了几年营销数据体系,总结下来营销自动化对OLAP的需求大概可以分成四类:
第一类:多维任意圈选。运营圈人永远不会只按一个维度筛选。"性别=女、城市=上海、最近一次购买在60天内、看过A商品、没领过B优惠券",这种多条件组合查询是日常操作。维度可能是用户属性、行为事件、订单记录,甚至可能是历史营销触达记录。这就要求OLAP引擎能高效处理大宽表上的条件过滤,并且支持Bitmap这类快速过滤索引,否则圈人的响应时间根本没法看。
第二类:行为序列分析。营销里特别关注用户在特定事件前后的行为路径。"先领券、再加购、最后下单"是一个转化路径,"加购后超过2小时未支付"是流失信号。这类需求需要对事件按时间排序、做窗口计算,天然适合用明细表加窗口函数来解决,但对引擎的解析能力和计算效率要求不低——尤其当事件表达到几十亿行时,一段复杂的LEAD/LAG窗口查询就可能让引擎跑几分钟。
第三类:精确去重与留存计算。连续几天的推送活动,到底覆盖了多少独立用户?这个月做过A活动又做过B活动的用户重叠了多少?日活、月活、留存这类指标全部依赖精确去重。ClickHouse的近似去重函数在早期版本误差挺大,做营销人群盘点时很难接受,而营销圈人又恰恰不能接受"约等于"——你给运营说"大概5万人",运营按5万人准备了预算,结果实际触达4.5万还好说,如果触达了5.5万,超发的成本就尴尬了。
第四类:高并发自助查询。运营团队不可能每次都提需求让数仓同学写SQL。自助分析工具接上OLAP引擎之后,查询发起量会上一个量级。引擎不仅要能跑得动复杂查询,还要能扛住几十个人同时点来点去。单机版解决方案在这里基本出局,这也是后来我们放弃了一些轻量级方案的原因。
这四类需求合在一起,基本排除了"用MySQL硬扛"和"用传统BI报表定时刷新"这两条路。
1.3 为什么不是所有OLAP都适合营销场景
这几年市面上OLAP引擎一大堆,但并不是每个都适合营销自动化。我见过有的团队一开始选了偏日志分析的引擎,结果做圈人时发现精确去重和数据更新怎么都搞不定;也见过团队选了事务能力很强的数据库,结果在亿级明细聚合上直接拉胯。
营销自动化的特点在于它"既要又要还要":既要有明细查询能力,能看单个用户的行为轨迹;又要有聚合分析能力,能看一群人漏斗转化;还要有实时更新能力,用户刚下完单,标签就要立刻变化。这三者叠加,市面上能同时做好的产品不多。我们在选型时最终圈定了StarRocks和Doris这一类的MPP分析数据库,具体原因后面在架构演进的部分细说。
2. 多源数据接入:先解决ID打通和口径统一,再谈架构
2.1 营销数据源全景盘点
做营销自动化的数据平台,第一件事不是选OLAP引擎,而是把数据源盘清楚。我见过太多团队一上来就急着搞引擎选型,结果引擎选完了,发现上游数据都没接进来,架构白搭。
营销自动化的数据源,按我的习惯分成五类:
第一类:线上行为数据。来自网站、App、小程序上的埋点事件:浏览、点击、加购、收藏、搜索、注册、登录。这类数据量最大、实时性要求最高,一般通过SDK采集后直接进Kafka。埋点数据的质量直接决定后面的分析能不能做,所以埋点规范从第一天就要定死。
第二类:交易与订单数据。订单、支付、退款、发货。这类数据通常在业务数据库里(MySQL或PG),量级比行为数据小,但准确性和实时性要求高。通过binlog订阅(Canal或Debezium)同步,或者业务方直接调用接口上报。注意:订单数据的金额字段一定要用分为单位存整数,避免浮点数精度问题——这是好多团队踩过的老坑。
第三类:CRM与客服数据。用户基本信息、会员等级、积分、客服沟通记录、投诉工单。这些数据更新不那么频繁,但维度丰富,是用户画像的重要来源。CRM数据的质量通常比较感人,手机号格式五花八门、同一个用户建档建了好几份,接入时要做好清洗。
第四类:广告投放数据。各广告平台的曝光、点击、消耗、转化回传。这类数据一般通过平台开放API拉取,频率从小时级到天级不等。麻烦的是每个平台的字段定义都不一样,曝光口径更是五花八门,同一个用户看完广告到底算不算"有效曝光",不同平台能给出三个答案。
第五类:触达记录数据。推送、短信、邮件、站内信的发送记录、送达状态、用户点击行为。这类数据是营销闭环的关键,不接进来就没法做触达效果分析,也不知道哪条推送把人给推烦了导致卸载。
把五类数据源拉通之后,才能真正开始设计OLAP架构。所有数据最终汇入Kafka,实时和离线双通道进入分析引擎,再由上层营销引擎和自助分析平台消费。
2.2 ID打通:匿名ID与实名身份的映射
多源数据合并时遇到的第一只拦路虎是ID统一。
用户没登录的时候,你只能拿到设备ID或者Cookie;登录之后你有了用户ID;他在小程序里授权给了你UnionID;下单时填了手机号;在门店办卡的时候录入了姓名和手机号。同一个人在不同时间、不同渠道留下的ID完全不一样。如果不打通,圈人就会出现误判——同一个用户在App里没登录的行为算了一个人,小程序里的行为又算了另一个人,最后推送重复发送,体验稀烂。
我们当时的做法是维护一张ID映射表,也就是常说的ID Graph。以用户ID为核心主键,把设备ID、UnionID、手机号等辅助标识全部挂在一个统一ID下面。每条新事件进来,先做一次ID解析,把匿名标识映射到统一ID上,然后再落到OLAP的事件表里。实时识别和离线回溯用两套逻辑:实时链路用Redis里的映射缓存,保证低延迟;离线链路跑Spark任务做全量ID合并,解决实时链路识别不了的老用户和长周期关联。
这个环节容易出问题的地方在于:一个手机号可能被多个用户使用过(比如情侣共用、二手设备),一个设备也可能被多个人登录过。完全精准的ID打通在现实里是不存在的,只能设定置信度规则。我们的原则是:涉及花钱的营销(比如发优惠券、返现),只用高置信度的打通结果;常规内容触达可以放宽条件。这个原则屡试不爽,本质上是把数据工程的容错思路用在了营销业务上。
2.3 口径统一:行为、时间、指标的三个战场
ID打通之后,第二个大坑是口径统一。这个比ID打通还要命,因为口径不统一,技术架构再漂亮,业务上也会互相掐架。
行为事件的定义口径。一个"加购"事件,是从提交订单页面触发的,还是点击"加入购物车"按钮就触发?同一个按钮,在H5和App两个端的埋点参数可能不一样。不同版本App之间,参数名可能都改过。这些都需要在数仓的DWD层做标准化,把五花八门的原始事件清洗成统一schema的"标准事件"。
时间口径。"某用户昨天加购"里的"昨天",是按用户本地时间算,还是按服务器时间算?运营在上海看数据和用户在纽约产生行为,如果不统一时区,做出来的圈选必然有偏差。我们统一用业务国家时区加UTC双字段存储,分析默认按业务方指定时区执行。
指标口径。GMV到底是支付成功的金额,还是下单未取消的金额?"领取优惠券人数"是去重后的用户数还是领取次数?这些指标口径必须收敛到一个指标平台上,所有人查同一个指标定义。否则就会出现运营说A方案效果好、市场说B方案效果好,最后发现两边连"效果"的定义都不一样。
口径统一这种事,靠技术手段只能做到一半,另一半靠组织和流程。我的经验是:数据团队牵头,拉上运营、产品、商业分析每个季度做一次口径评审,把指标字典维护成公司级资产。这活儿不性感,但不做的话,后面OLAP里存的数据再准,业务方也不敢信。
2.4 数据质量:重复、乱序、延迟的应对
最后聊一下多源接入的数据质量问题。这玩意儿不解决,OLAP只能存一堆垃圾。
重复数据。埋点SDK重试上报、API拉取重复、binlog传递重复消费,都会导致重复事件。处理方式是在OLAP层对事件ID做去重——如果引擎支持主键模型或者UPSERT,就设置唯一键;同时计算侧做幂等设计,保证重复消费同一批数据最终结果一致。
乱序数据。用户手机断网之后恢复,一批事件延迟上报,事件产生时间和到达时间可能差了半个小时甚至一天。实时链路里要容忍乱序,Flink作业要设置Watermark策略。OLAP引擎里一般以事件发生时间为准做窗口计算,而不是以入库时间为准。
延迟数据。离线任务跑完了,实时表里又补进来一段历史数据,导致两个链路的数据对不上。我们的策略是:实时表以"当日实时+昨日及以前由离线回刷"的方式修正,每天凌晨跑一次对账任务,发现差异自动触发数据订正。这套机制保障了OLAP里数据的可解释性。
3. 架构演进实录:从MySQL硬查、离线数仓到实时OLAP
3.1 第一阶段:MySQL直查报表体系
我接手营销数据平台的时候,架构非常简单粗暴:业务库MySQL,再加一堆定时脚本,每天早上从业务库抽数,生成几十张汇总报表,扔到报表平台里。运营要看数据,就去看报表;要做人群,提需求给数仓,数仓写SQL导出来给运营。
这个阶段的特点是一个字:慢。报表T+1,运营看到的全是昨天的数据。真正要圈人的时候,数仓临时写一个SQL在MySQL上跑,跑个十几分钟超时是常事。更要命的是,MySQL上跑重型查询会直接影响线上交易业务的性能。有一次运营在大促前拉人群做优惠券推送,一条复杂查询把主库的慢查询日志刷了好几屏,DBA差点掀桌子。
现在回头看,第一阶段最核心的问题不是"没有OLAP",而是整个数据链路没有和业务解耦。查询和分析负载压在生产库上,安全性和稳定性都堪忧。这个阶段其实谈不上架构,只是能用而已。KPI是"能出数",不管是跑半个小时还是一个小时,只要最终能导出一份人群名单,运营就谢天谢地了。
3.2 第二阶段:离线数仓 + 独立OLAP引擎
被MySQL折磨了小半年之后,我们做了第一次架构升级。核心思路是:把数据从业务库复制出来,建立独立的数仓,然后用一个新的OLAP引擎承接分析查询。
具体链路是:业务库binlog加埋点数据进Kafka,同步任务写入Hive构建离线数仓(ODS原始层、DWD明细层、DWS汇总层、ADS应用层),再用调度任务把结果同步到OLAP引擎里。运营的查询全部打在OLAP引擎上,不再碰业务MySQL。
这个阶段解决了两个大问题:一是把分析负载和生产库隔离了,DBA终于不找我了;二是数仓分层之后,口径统一有了落地的载体,不同报表终于能对上一个数了。
新问题也冒出来:T+1的时效性仍然没有破。运营的圈人需求还是只能圈"截止到昨天的用户",对于"今天刚加购但没付款"这种实时流失用户完全无能为力。而营销恰恰是最吃"实时"的场景——用户在犹豫的当下推一张券,和第二天再推,转化率天差地别。我们测过一组数据:加购后30分钟内发券的转化率,是次日发券的3倍左右。这个数字直接推动了实时链路的建设。
3.3 第三阶段:实时与离线双通道融合
于是我们进入到第三个阶段:数据双通道架构。实时链路用Kafka接埋点和业务binlog,Flink做实时ETL,结果直接写入OLAP引擎;离线链路继续跑Hive数仓,每天定时回刷。这样OLAP引擎里既有实时更新的数据,又有全量历史数据,圈人和分析可以同时基于一份数据源。
这一阶段我们还做了一件关键事:把数据服务层拎出来。OLAP引擎不直接暴露给业务系统和自助分析工具,而是通过一个统一的数据服务API层对外提供服务。营销引擎的圈人接口、人群计算任务、运营的自助分析平台,都走这层。好处是:OLAP引擎的运维变更不影响上层业务,权限控制也集中在一个地方。
经历过这三个阶段之后,我个人的体会是:架构演进不是一蹴而就的,每一次升级都是被业务痛点推着走的。第一阶段的痛点是不能查询,第二阶段的痛点是查询太慢,第三阶段的痛点是不实时。每一个阶段的技术选型都要匹配当下的业务复杂度和团队能力,上来就上最重的架构大概率会翻车。
3.4 选型对比:StarRocks、Doris、ClickHouse怎么选
说到OLAP引擎选型,这是我被问得最多的问题。我们的场景经过评估,最终选择的是StarRocks。结合使用体验,我说下三个引擎的差异。
| 能力维度 | ClickHouse | Doris | StarRocks |
|---|---|---|---|
| 单表查询性能 | 很强 | 强 | 强 |
| 多表JOIN | 相对弱 | 好 | 好 |
| 精确去重 | 弱于专业OLAP | 好 | 好 |
| 数据更新(UPSERT) | 弱,需要Merge | 好(主键模型) | 好(主键模型) |
| 高并发点查 | 弱 | 好 | 好 |
| MySQL协议兼容 | 部分 | 是 | 是 |
| 实时写入能力 | 强 | 强 | 强 |
ClickHouse的优势是单表查询性能极强,导入速度也快,生态成熟,社区活跃。但它在高并发查询、多表JOIN、精确去重、数据更新这几个能力上相对偏弱。早期版本做数据更新要Merge函数,实时数据修正非常痛苦。如果你的场景主要是日志分析、可观测性或者固定报表,ClickHouse完全够用。
Apache Doris和StarRocks,同属于MPP分析型数据库,核心优势是支持高并发点查、支持主键模型做实时更新、JOIN性能好、兼容MySQL协议。这两个引擎对有"既要又要还要"需求的营销场景特别合适:既要实时更新用户状态,又要跑复杂的圈人查询,还要支撑高并发的自助分析。
在Doris和StarRocks之间怎么选,我的看法是:如果团队有人能深度维护开源Doris代码,或者有相关经验,选Doris没毛病;如果希望有商业资源支持、版本迭代更快,想少踩一些版本兼容的坑,可以选StarRocks。我们选StarRocks的主要原因是营销场景对数据更新的需求太频繁,而它在主键模型上的成熟度当时更胜一筹。
选型不能只看引擎能力,还要看团队的技术储备。一个只有MySQL经验的小团队,上StarRocks也要学一段时间;但如果连基础的数仓分层都没理清楚,贸然上复杂架构迟早要还债。我的建议是:先搞清楚业务最核心的一两个痛点,再选一个能解决痛点且学习成本可接受的引擎,千万别奔着"技术先进"去选。
4. 营销场景的核心数据模型设计
4.1 分层建模:DWD明细、DWS汇总、ADS应用
传统数仓的分层方法论放在OLAP架构里依然成立,但具体表的设计思路跟Hive数仓有差别。OLAP引擎的数据存储和计算资源都比HDFS贵得多,不可能像Hive那样把所有明细原样铺一层。所以在OLAP里的分层要更"抠门"。
我们的做法是三层:
- DWD明细层(实时+离线合并):核心是用户行为事件表,只保留营销分析真正用得到的字段,JSON属性里高频查询的指标字段单独拆列存储。事件表是OLAP里最大的一张表,也是圈人查询最主要的数据源。
- DWS汇总层:按用户、按天计算的轻度汇总。比如用户每日行为汇总表(每个用户每天的浏览次数、加购次数、下单次数、金额)、用户活跃状态表(最近一次活跃时间、连续活跃天数)。
- ADS应用层:面向具体业务场景的宽表和预计算结果。比如人群圈选结果表、营销活动效果明细表、漏斗分析结果表。
分层的好处是:DWD解决"能不能查",DWS和ADS解决"查得快不快"。日常运营的圈人查询,80%可以直接命中DWS和ADS,只有深度的行为序列分析才需要下探到DWD。
4.2 用户事件表与用户宽表
事件表是整个OLAP模型里最重要的表。我们的设计大概是这样的:
CREATE TABLE dwd_user_event ( event_id VARCHAR(64), user_id BIGINT, device_id VARCHAR(64), event_type VARCHAR(32), -- view/add_to_cart/payment... event_time DATETIME, -- 事件发生时间(业务时区) event_time_utc DATETIME, p_date DATE, -- 事件发生日期分区字段 session_id VARCHAR(64), page_id VARCHAR(64), product_id VARCHAR(64), order_id VARCHAR(64), channel VARCHAR(32), -- 渠道标识 extra_props JSON, ... )事件表的设计有几个要点:用event_id做主键,支持幂等写入,重复消息进来自动去重;事件类型单独一列,不要混在JSON里,方便过滤和索引;高频查询的字段(商品、订单、渠道)显式拆列,低频字段放JSON;用事件发生时间分区,而不是入库时间分区,否则离线回刷和实时对不上。
用户宽表则是把用户的静态属性和动态状态合并成一行。静态属性来自CRM和注册信息,动态状态比如"最近一次购买时间""累计消费金额""当前会员等级"。"近30天活跃天数"这种频繁变化的指标,如果放宽表会导致频繁更新,我们一般把它放在DWS的行为汇总表里,通过JOIN拿。
宽表设计的时候我踩过一个坑:一开始把所有指标都往一张表上堆,字段堆到两百多个,导入性能直线下降,而且一张表里指标更新周期各不相同,主键模型每天UPSERT大几百GB数据,慢得没法看。后来把宽表拆成"用户维度表"(低频更新)和"用户行为指标表"(高频更新)两张,JOIN查询的性能反而更好了。所以在OLAP里,宽表不是越宽越好,更新频率相近的字段才适合放一起。
4.3 漏斗分析与行为路径建模
营销分析绕不开漏斗。从曝光到点击、从加购到支付,每一步有多少人、流失在哪一步,是运营每天都要看的数据。
漏斗分析在OLAP里有个经典的写法:用窗口函数给每个用户的行为路径按时间排序,然后判断每一步是否发生在指定窗口内。举个例子,一个四步漏斗:浏览商品、加购、提交订单、支付成功,要求60分钟内完成。
WITH user_journey AS ( SELECT user_id, event_type, event_time, LEAD(event_type, 1) OVER (PARTITION BY user_id ORDER BY event_time) AS next_event, LEAD(event_time, 1) OVER (PARTITION BY user_id ORDER BY event_time) AS next_time FROM dwd_user_event WHERE p_date >= :start_date AND p_date <= :end_date ) SELECT COUNT(DISTINCT CASE WHEN event_type = 'view_product' THEN user_id END) AS step1_users, COUNT(DISTINCT CASE WHEN event_type = 'add_to_cart' AND next_event = 'create_order' AND next_time <= event_time + INTERVAL 60 MINUTE THEN user_id END) AS step2_users, ... FROM user_journey这段SQL的漏斗人数是逐步收敛的,因为条件里用到了下一步事件。实际查询中如果漏斗步骤多达七八步,需要仔细设计SQL,避免对事件表做多次全表扫描。我们的经验是把漏斗定义提前配置好,用物化视图或者异步任务把常用漏斗的结果预计算到ADS层,查询时直接取结果。
行为路径分析比漏斗更复杂一点,它要还原每个用户在某段时间内的行为序列。比如要看"加购未支付用户在下单前的浏览路径",需要把事件明细按用户分组、按时间排序后截取序列。OLAP引擎对这类查询一般都能支持,但数据量大时开销不小,建议只跑近30天以内、且限定用户群的查询。
4.4 归因模型:多触点归因在OLAP里的落地
营销效果归因是另一个绕不开的话题。用户可能先在朋友圈看到广告,然后在抖音刷到信息流,再被社群分享文章种草,最后通过搜索进入网站下单。这一串路径里,每个触点各自贡献了多少转化?归因模型就是回答这个问题的。
归因模型有很多种:末次点击归因、首次点击归因、线性归因(每个触点均摊)、时间衰减归因(越接近转化的触点权重越高)、位置归因(首位和末位各占40%,中间均摊20%)。在OLAP里落地归因计算,核心是能把一个用户的全部触点按时间排序,然后根据不同模型给每个触点分配权重。
我们用的是明细事件表加统一归因SQL的方式。每个触达事件关联一个曝光时间,每笔订单关联一个下单时间,然后在用户维度上做时间序列对齐:
SELECT t.touch_channel, SUM(CASE WHEN :attribution_model = 'last_click' THEN CASE WHEN t.is_last THEN t.order_amount ELSE 0 END END) AS attributed_gmv FROM ( SELECT a.user_id, a.touch_channel, a.touch_time, o.order_amount, ROW_NUMBER() OVER (PARTITION BY a.user_id, o.order_id ORDER BY a.touch_time DESC) AS rn, (MAX(a.touch_time) OVER (PARTITION BY a.user_id, o.order_id)) = a.touch_time AS is_last FROM dwd_touch_event a JOIN dwd_order o ON o.user_id = a.user_id AND o.order_time BETWEEN a.touch_time AND a.touch_time + INTERVAL 7 DAY ) t WHERE t.rn <= 10 GROUP BY t.touch_channel;归因计算的查询通常跑得比较重,特别是把7天回溯窗口内所有触达和所有订单做匹配的时候。我们是每天凌晨跑一次T+1全量归因,写入ADS结果表,供报表和投放优化使用。实时归因只在探索性项目里做过,因为实时数据波动大,直接用来做投放调整容易产生误判。
5. 实时圈人与自动触达:OLAP怎么和营销引擎配合
5.1 动态人群圈选的技术实现
如果说前面讲的都是"能查",那营销自动化最核心的应用就是"能用查出来的结果去触发动作"。这里的关键是圈人任务的调度和执行。
圈人任务分两类:
离线圈人(批量):运营配置条件,调度系统每天跑一次,结果写入人群表。适合周期性营销,比如"每周一给上周活跃但未下单用户发券"。
实时圈人(动态):用户一发生某个行为,立刻判断是否满足触发条件,满足就进人群并触发触达。比如用户加购后30分钟未支付,系统马上推送一张优惠券。这类场景对OLAP查询的实时性要求很高。
实时圈人的实现我们用了两种方案。第一种是纯流式计算:Flink消费行为事件流,维护用户在状态中的画像信息,命中规则就输出触达指令。这个方案实时性最好,毫秒级响应,但规则都写在Flink作业里,运营不能自助配置,灵活性很差。
第二种方案是"流批一体":Flink作业只负责把实时事件写入OLAP引擎,触发判断由营销引擎实时发起查询。营销引擎收到"加购"事件后,查一次OLAP——看这个用户是否满足其他触发条件(比如近30天活跃次数、是否已领券、是否黑名单)。这个方案的好处是判断条件可以灵活配置在营销引擎里,代价是每次事件都要发一次OLAP查询,对引擎的高并发点查能力要求很高。
选型上,我们最终是两种方案混合用的。高频且规则固定的场景(比如加购未支付)用纯流式计算,低频且规则复杂多变的场景用OLAP实时查询。混合用最大的好处是成本和灵活性之间的平衡。纯流式计算省掉了每次查询的开销,但新增一个规则就要改作业发版,这在运营跑活动的时候时效性跟不上。
5.2 触达闭环与效果回流
圈人只是营销自动化的起点,真正的闭环是触达和效果回流。用户被圈选出来之后,要通过推送、短信、优惠券等渠道触达。触达完,还要把触达结果、用户后续行为回收进OLAP,形成闭环。
整个闭环的数据链路大概是这样的:
- 营销引擎查询OLAP,拿到人群用户ID列表;
- 通过数据服务层调用触达平台接口,发送推送、短信、发放优惠券;
- 触达平台把发送状态、送达状态、点击记录写回Kafka;
- Flink消费Kafka数据,更新OLAP里的触达明细表;
- 后续用户行为事件照常进入事件表,营销效果分析时,把触达记录、行为事件、订单数据拉通分析。
这个闭环里最容易出问题的是"触达记录和行为事件的时间顺序"。比如用户先收到推送,2小时之后下单,这个订单算不算推送带来的转化?如果触达记录延迟入湖,对账就会出问题。我们的做法是触达事件进入OLAP时,除了写入时间外,必须带上触达发送的原始时间戳,所有效果分析一律按原始时间戳做匹配,不依赖入库顺序。这一点写进了数据规范,谁违反谁负责,因为一旦依赖入库顺序,数据回刷或者断点续传的时候,归因结果就会乱掉。
5.3 时效性分级:不同营销场景对延时的不同要求
营销自动化的延时要求其实不是统一的。做了一段时间后,我发现可以把营销场景按对延时的敏感度分成三档:
实时档(秒级~分钟级):加购未支付挽回、价格变动通知、支付失败提醒。用户在当下最有行动意愿,晚推一分钟都可能流失,必须用流式计算或者高频查询。这类场景对OLAP的高并发点查和Flink的状态管理能力要求最高。
准实时档(分钟级~小时级):文章发布后给高活跃用户推送、活动开始前的预告提醒。对分钟级延迟不敏感,用Flink微批或者每5分钟一次的批量任务就可以覆盖。这类场景占了营销自动化的很大一部分,也是ROI提升最明显的部分。
离线档(天级):周报、月报、人群包复盘、周期性营销任务。"昨天数据"完全够用,走离线调度最省成本。
建议做架构设计时把时效性分级画成一张矩阵图,跟业务方对齐清楚:哪些场景可以接受延迟、哪些不行,然后针对每一档设计不同的链路,避免一刀切全都走实时链路。实时链路的成本和运维复杂度都比离线高不少,如果所有场景都上实时,数据团队会被运营的各种临时需求搞到崩溃。
6. 查询调优实战:营销场景下OLAP优化的踩坑记录
6.1 排序键与分区设计踩过的坑
OLAP引擎性能好不好,排序键和分区键的设计占了很大比重。这块我们踩过的坑,可以写一篇单独的长文。最典型的一个:一开始我们把事件表的分区粒度设计成小时,排序键用的是event_id,结果查询全表扫描严重,因为运营的查询基本都是按时间范围和用户维度过滤的,很少直接按event_id查。
后来我们把分区改成按天,排序键调整为(event_time, user_id, event_type),效果立竿见影。大多数圈人查询都能通过分区裁剪和前缀过滤把扫描范围缩小到很小,查询耗时从几秒钟降到了几百毫秒。这个优化没有改任何业务代码,纯靠调整表结构就实现了,所以表结构设计在OLAP里是性价比最高的优化手段。
在设计排序键的时候,我的建议是:把查询中频率最高的等值条件和范围条件放在最前面,而不是按字段顺序排。像"渠道=APP"这种等值条件、"时间范围"这种范围条件,放在前面能最大程度发挥引擎的索引剪枝能力。如果一张表同时被多种查询模式使用,排序键只能兼顾少数几种,那就需要按查询模式拆表。比如分析常用的维度排序键和圈人常用的时间排序键,就分别建了两张表,虽然有点冗余,但换来了查询的稳定和快速。
6.2 数据倾斜:高价值用户的查询反噬
数据倾斜在营销场景里特别明显。头部用户贡献了大部分GMV,他们的行为事件量也远高于普通用户。一个爆款商品的浏览事件可能集中在同一天同一个商品ID上,查询这个商品的漏斗时,计算节点负载极其不均,慢节点拖垮整个查询。
我们的处理办法是对热点数据进行打散。比如在事件表里额外增加一个分桶键,把热点用户的事件数据分布到更多分片上;或者把热点商品的查询改写为按商品维度预计算的ADS表,绕开明细聚合。另外,对常见的"爆款分析""热门活动分析"这种高访问量指标,直接跑定时物化视图,查询直接走预计算结果,基本能消除数据倾斜带来的影响。
需要提醒的是,数据倾斜不是上线时才发现的,通常在压测阶段就能暴露出来。我们的经验是:新表上线前,用线上真实数据分布做一轮pre-check,专门看高频用户、高频商品、高频渠道这几个维度的数据占比。如果某个单一维值的数据量超过全表的10%,就要提前做打散策略,不要等生产环境挂了再补救。
6.3 实时与离线数据一致性核对方法
双通道架构下,实时表和离线表的数据一致性是个长期工程。我们建立了一套每日对账机制:按天校验事件表里各事件类型的总量差异;按用户维度抽检实时表和离线表在部分用户上的行为序列差异;对核心指标(如GMV、加购人数)设置差异告警阈值,超过阈值自动触发离线表回刷实时表对应分区。
这套对账机制跑下来,我们发现大多数数据不一致都源于实时链路的乱序和迟到数据。在Flink里适当放宽Watermark延迟阈值之后,实时和离线的对账通过率从90%提到了99%以上。剩下不到1%的差异,人工复核基本上都是埋点定义的问题,不是OLAP引擎的问题。
对账的另一个作用是反向推动埋点质量。有些埋点在发版时写错了参数,实时链路和离线链路解析的结果不一样,对账任务就会揪出来。刚开始运营很反感这类报错,觉得影响他们看数,后来发现埋点错一次导致的错误决策成本远高于修正成本,也就接受了这种"被迫严谨"。
6.4 缓存与数据服务层的兜底策略
最后说下OLAP高并发查询的兜底。营销引擎在某些场景下(比如大促时集中圈人)会在短时间内发起大量查询,OLAP引擎再强也扛不住峰值流量直接怼。我们在大促前做了一些保护措施:
一是数据服务层加了一层Redis缓存,把运营经常用的固定人群结果缓存到Redis,查询能命中缓存就直接返回,不压OLAP。二是对预估执行时间超过阈值的查询做排队处理,避免重型查询把引擎的资源占满,导致核心查询饿死。三是在大促期间关闭一些非核心的深度分析类查询入口,优先保障圈人链路稳定。
这套兜底策略本质上是在"数据新鲜度"和"系统稳定性"之间做权衡。营销场景里,一个查询慢几秒钟不会死人,但整个系统挂了就是事故。所以我在设计数据服务层时,优先保证系统的可用性和可降级性。缓存失效策略用的是"缓存5分钟+主动失效"双机制,既保证数据的相对新鲜,又能在运营改动活动配置后立刻刷新人群结果。
说到最后,我还想分享一个从多次事故里总结出来的经验:很多人觉得OLAP调优就是把引擎参数调好、SQL写漂亮,但在真实的营销自动化场景里,数据链路任何一个环节出了问题,最后都表现为"查询变慢"或者"数不对"。所以排查问题的时候,不要只盯着OLAP引擎本身,要从数据源到调度到服务层全链路看一遍。多源数据的OLAP架构演进,本质上不是选一个引擎那么简单,而是把数据采集、ID打通、口径统一、模型设计、实时计算、查询调优串成一条完整的链路。每个环节都有坑,但每个坑踩过去之后,整套系统就会变扎实不少。希望我这几年的经历,能让你在搭自己这套体系的时候少走几步弯路。