从用户画像到推送频控:大数据推荐与推送系统架构实战

发布时间:2026/9/8 11:10:16
从用户画像到推送频控:大数据推荐与推送系统架构实战 “懂用户”和“惹人烦”之间往往只差一次推送的距离。“大数据这么会推那就多推”这句话放到工程语境里其实是在问一个问题当推荐策略跑通之后系统能不能支撑更大量的推送、更精准的触达、更复杂的业务规则这次我们不聊概念直接拆一套可落地的大数据推荐与推送系统参考架构。这套架构覆盖用户画像、人群圈选、召回精排、推送频控、接口 API、批量任务和集群部署。无论你是准备大数据毕业设计还是在公司里接手用户增长相关的大数据开发任务都可以按这条链路走。本文会讲清楚一条完整的“数据采集 — 画像加工 — 推荐策略 — 消息触达 — 效果回流”链路长什么样每个环节需要哪些组件、哪些表、哪些接口批量推送任务怎么设计才不会把服务打挂以及最容易踩的坑有哪些。建议收藏备用。1. 核心能力速览能力项说明项目定位大数据推荐与推送系统参考架构侧重数据链路、人群策略、推送服务和批量任务数据链路埋点日志 → 消息队列 → 实时/离线计算 → 用户画像库 → 推荐服务 → 推送任务表 → 触达渠道推荐策略规则召回、热门召回、协同过滤、CTR 预估排序、多样性重排推送能力用户级单推、人群批量推送、定时推送、频控/疲劳度控制、退订过滤部署形态Hadoop 集群 / 伪分布式Flink、Spark、Kafka 可按需组合接口能力预留推送任务提交接口、人群圈选接口、效果回传接口批量任务支持离线人群文件批量导入、任务状态机管理、失败重试合规边界用户授权、数据脱敏、退订优先、个人信息保护要求适合场景大数据毕业设计、实训项目、企业内部用户触达中台、推荐系统入门实践这套架构的定位是“参考实现”。实际落地时组件选型、表结构、接口字段都要按自己的业务场景调整但整体分层思路是可复用的。2. 适用场景与使用边界先说清楚这套东西适合谁。第一类是正在准备大数据毕业设计的同学。标题“大数据这么会推那就多推”很适合作为一个演示项目的出发点做一个“用户推荐与消息推送系统”把用户画像、实时计算、推荐策略、推送服务全部串起来答辩时可以完整讲出技术链路。第二类是在公司里做大数据开发、用户增长或私域运营的工程师。很多公司已经接入了神策、友盟或自研埋点数据有了但缺少一套“从画像到触达”的工程化方案。这套架构可以作为中台设计参考。第三类是准备大数据面试的人。面试官问“你做过推荐系统吗”时能说清楚召回、粗排、精排、重排的职责能画出数据链路图能讲明白批量推送怎么做幂等和重试会是非常加分的技术点。再说边界。这套架构解决的是“推送效率”和“系统稳定性”问题不能解决“内容本身不行”的问题。如果推荐内容质量差推送越多用户退订越快。同时必须强调合规。做用户画像和推送之前要确认数据来源有用户授权涉及手机号、设备号、行为轨迹等个人信息时必须做脱敏和权限管控。每一次推送都要给用户提供退订入口退订记录要立刻生效不能出现“用户退订了还继续收到消息”的情况。3. 总体架构与数据流整套系统的核心是“一条闭环数据流”用户发生行为行为被采集画像被更新策略算出内容推送触达用户用户再次行为回流。3.1 分层结构可以把系统分成五层层次职责常用组件数据采集层客户端埋点、业务日志采集、数据库 Binlog 采集Flume、Logstash、Canal、Kafka计算存储层实时计算、离线计算、画像存储Flink、Spark、Hive、HBase、ClickHouse、Redis策略层用户分群、推荐召回、排序、画像查询Spark Job、ES、HBase、Redis触达层推送调度、频控、任务状态管理自研推送服务、App 推送、短信、站内信效果层曝光、点击、转化回传AB 实验分析Kafka、Hive、OLAP3.2 数据流示例一条典型的推送链路是这样走的用户行为日志 → Kafka(click_topic) ├─ Flink 实时消费 → Redis 在线特征 └─ Hive/Spark 离线加工 → 用户标签宽表 ↓ 人群圈选服务 ↓ 推荐召回 精排 ↓ 推送候选队列 ↓ 推送调度 频控 ↓ App Push / 短信 / 站内信 ↓ 曝光/点击回传 → Kafka(back_topic)这条链路的重点在于“回流”。如果只做推送不做效果回传推荐策略就没办法迭代等于闭着眼睛投放。3.3 数据格式示例用户行为日志统一使用 JSON 格式写入 Kafka字段不能太随意否则下游解析成本很高。{ event_id: uuid-001, user_id: u_10001, event_type: item_click, item_id: item_2034, scene: homepage, device_id: device_xxx, platform: android, ts: 1719216000000, extra: { from_source: push, push_task_id: task_2024062001 } }push_task_id这种字段建议从第一天就加进来。否则推送触达之后后台根本没法统计某次推送带来了多少点击转化。4. 大数据集群部署与环境准备4.1 集群部署策略学习环境、毕业设计环境和企业生产环境的部署策略完全不同。环境类型部署方式说明本地学习/毕业设计单机伪分布式一台 16G 内存以上的机器部署 Hadoop、Hive、Spark、Kafka团队测试环境3 台虚拟机/云主机NameNode DataNode 分离YARN 资源队列隔离企业生产环境多节点集群 云托管高可用、数据副本、混合部署或存算分离如果是毕业设计不建议一开始就搭 5 台机器的集群。单机伪分布式能跑通完整链路已经足够展示工程能力。4.2 组件规划参考组件用途备注Hadoop HDFS离线数据存储伪分布式即可Hive离线标签加工建数仓分层表Spark人群圈选、离线推荐SparkSQL 为主Kafka日志采集、消息解耦至少 3 个 topicFlink实时画像更新选修能加分HBase/ClickHouse画像在线查询HBase 适合 kvClickHouse 适合分析Redis缓存、频控计数必须MySQL推送任务、用户退订表必须4.3 启动命令模板不同的发行版启动命令有差异下面给的是通用模板。# 启动 HDFS 与 YARNHadoop 发行版不同路径不同 $HADOOP_HOME/sbin/start-dfs.sh $HADOOP_HOME/sbin/start-yarn.sh # 启动 Kafka按实际安装目录调整 bin/kafka-server-start.sh config/server.properties # 提交 Spark 人群圈选任务 spark-submit \ --master yarn \ --deploy-mode client \ --executor-memory 2g \ --num-executors 4 \ job/crowd_select.py启动脚本不要写死在代码里建议用一个env.sh统一管理组件路径和端口。第一次跑通环境后把整个启动流程保存为一条脚本后续重置环境会非常省事。5. 用户标签与人群圈选推送的第一步不是发消息而是圈出“这批消息该发给谁”。5.1 用户标签宽表设计用户画像数据最终要落成一张宽表每个用户一行行为标签存在 Map 字段里。CREATE TABLE user_profile_daily ( user_id STRING COMMENT 用户ID, tags MAPSTRING, STRING COMMENT 标签集合: 如 age_band20_30, is_new1, active_days_7 INT COMMENT 近7天活跃天数, order_cnt_30 INT COMMENT 近30天订单数, last_active_ts BIGINT COMMENT 最后活跃时间戳, opt_out_push INT COMMENT 是否退订推送: 0否 1是, register_date STRING COMMENT 注册日期, etl_date STRING COMMENT 数据分区日期 ) PARTITIONED BY (etl_date STRING);注意opt_out_push必须在画像表里出现。很多推送事故都是因为没有过滤退订用户。5.2 离线人群圈选脚本以“近 7 天活跃、近 30 天未下单”的用户为例使用 SparkSQL 圈选。INSERT OVERWRITE TABLE push_crowd_daily PARTITION (etl_date ${target_date}) SELECT a.user_id, recent_active_not_buy AS crowd_type FROM user_profile_daily a WHERE a.etl_date ${target_date} AND a.opt_out_push 0 AND a.active_days_7 3 AND a.order_cnt_30 0这段 SQL 的核心是“过滤退订用户”和“用行为数据反推用户意图”。圈选规则不是拍脑袋定的建议每一类人群都对应一个业务目标比如“促进新客首购”“召回流失用户”“沉默用户激活”。5.3 实时人群更新离线画像一般 T1 更新但用户点击商品后希望立即收到推送时就需要实时链路。Flink 消费 Kafka 埋点日志实时更新 Redis 中的在线标签。// 伪代码示例Flink 消费点击事件更新 Redis 用户最近浏览商品 DataStreamClickEvent stream env.addSource(kafkaSource(click_topic)); stream.keyBy(ClickEvent::getUserId) .process(new RedisSinkFunction()) .name(update_user_recent_view);实时链路不建议承担太复杂的计算复杂规则尽量放到离线任务里。在线服务只做“最近一次浏览”“当前偏好品类”等轻量更新保证延迟可控。6. 推荐策略召回、粗排、精排、重排推荐系统不是“一个模型走天下”。在工程上推荐系统通常拆成召回、粗排、精排、重排四个阶段。这个结构在大数据面试里也是高频考点。6.1 召回层召回决定候选池的大小。召回太少容易导致内容单一召回太多会对后续排序造成压力。常见召回策略包括热门召回适合冷启动用户基于用户标签的规则召回例如最近 7 天浏览品类关联商品协同过滤召回基于用户行为相似度或者物品相似度向量召回在相对大规模的场景中使用 embedding 检索不是入门必需。毕业设计或者小型项目建议先用“规则召回 热门召回”组合链路简单效果可控。6.2 精排层精排负责把召回的候选物品排序。经典做法是训练一个 CVR 或 CTR 模型特征包括用户特征、物品特征、上下文特征、交叉特征。精排在离线阶段可以用 Spark 做特征拼接模型训练可以使用传统机器学习库比如逻辑回归或者 GBDT如果特征规模扩大到深度模型再考虑引入 GPU 训练但这不是必须的。6.3 重排层重排解决的是“排序分数最高不代表用户体验最好”的问题。重排要考虑业务规则例如同一个品类不能连续出现 3 个商品广告位要限流推送时每条任务不能含有重复内容。# 重排规则配置参考 rerank: max_same_category: 2 # 同一品类最多连续出现次数 max_ad_per_batch: 1 # 每批推送中广告位数量 exclude_read_items: true # 排除用户已读/已购内容6.4 推荐结果输出推荐服务最终输出的是一个有序内容列表并带上推送任务标识。这样后续做效果分析时能定位到某次推送到底推荐了哪些内容。{ user_id: u_10001, task_id: task_2024062001, items: [ {item_id: item_2034, score: 0.91, reason: user_recent_view}, {item_id: item_1052, score: 0.87, reason: hot_item} ] }很多项目做到这一步就停了但实际上“推送”才是业务价值的放大器。7. 推送服务与频控设计推送服务和推荐服务是两件事。推荐解决“发什么”推送解决“什么时候发、发给谁、发多少、怎么发不会打扰用户”。7.1 推送任务表设计所有推送动作先落地为任务而不是直接调用第三方推送接口。CREATE TABLE push_task ( id BIGINT AUTO_INCREMENT PRIMARY KEY, task_name STRING NOT NULL, task_type STRING COMMENT single:单用户, batch:批量人群, crowd_id STRING COMMENT 人群ID, push_channel STRING COMMENT app_push/sms/station_letter, content STRING COMMENT 推送内容模板, priority INT COMMENT 优先级: 1-10, 数字越大优先级越高, status STRING COMMENT pending/running/success/failed, total_count INT COMMENT 目标用户数, success_count INT COMMENT 成功推送数, fail_count INT COMMENT 失败推送数, request_id STRING COMMENT 幂等ID, create_time TIMESTAMP, update_time TIMESTAMP );任务表是批量推送的“中央调度器”。所有模块都围绕这张表运转调度模块扫描 pending 任务执行模块推送回传模块更新状态。7.2 频控与疲劳度控制“多推”不等于“滥推”。频控是整个系统最重要的底线设计。# 频控配置参考 push_frequency_control: global: max_push_per_user_per_day: 3 max_push_per_user_per_week: 10 min_interval_minutes: 120 min_interval_minutes_by_channel: app_push: 120 sms: 1440 station_letter: 30 blacklist: - user_blacklist - user_opt_out频控判断逻辑建议用 Redis 实现。每次推送前检查用户当日推送次数、最近一次推送时间、是否退订全部通过才允许写入发送队列。7.3 推送执行伪代码推送服务本质上是一个消费者任务从任务队列取用户逐个检查频控和退订状态然后调用第三方推送渠道。def push_to_user(user_id, content, channel): # 1. 检查退订 if is_opt_out(user_id): log_skip(user_id, opt_out) return {status: skipped, reason: opt_out} # 2. 检查当日推送次数 today_count redis_client.get(fpush_count:{user_id}:{today}) if today_count MAX_PUSH_PER_DAY: log_skip(user_id, frequency_limit) return {status: skipped, reason: frequency_limit} # 3. 检查最短推送间隔 last_ts redis_client.get(flast_push_ts:{user_id}) if last_ts and (now - int(last_ts)) MIN_INTERVAL: log_skip(user_id, interval_limit) return {status: skipped, reason: interval_limit} # 4. 调用第三方推送渠道 push_result channel_client.push(user_id, content) # 5. 更新计数 redis_client.incr(fpush_count:{user_id}:{today}) redis_client.set(flast_push_ts:{user_id}, now, ex86400) return push_result这段逻辑看起来很基础但大多数推送事故都出在“跳过了第 1 步”或者“第 2 步的 key 设计错误”上。频控 key 一定要注意过期时间否则用户天天收到推送的频率限制会失去效果。8. 接口 API、批量任务与 AB 实验8.1 推送 API 设计推荐系统和推送服务都建议提供 HTTP API方便上游业务方接入。// 请求 POST /api/v1/push/single { user_id: u_10001, content: { title: 你关注的商品降价了, body: 热门商品限时优惠点击查看, deep_link: app://product/2034 }, channel: app_push, priority: 5, request_id: req_202406200101 }// 响应 { code: 0, message: success, data: { task_id: task_20240620010001, status: accepted } }注意请求里必须带上request_id。否则上游业务方重试时会产生重复推送。Python 调用示例import requests url http://127.0.0.1:8080/api/v1/push/batch payload { crowd_id: crowd_recent_active_not_buy, content: { title: 新人专享优惠, body: 领取专属福利下单立减 }, channel: app_push, priority: 5, request_id: req_202406200102 } response requests.post(url, jsonpayload, timeout30) print(response.json())8.2 批量任务设计批量推送的核心是异步处理。接口提交人群之后不应该同步等待所有用户推送完成而是立即返回一个task_id后台 worker 从人群表分批读取用户执行推送。批量任务状态机pending(待执行) ↓ running(执行中) ←→ paused(暂停) ↓ success(全部成功) / partial_success(部分失败) / failed(失败)批量推送关键点分页读取用户避免一次性把几百万用户加载到内存每批次推送量做限制防止第三方渠道限流单条失败不中断整个批次记录失败原因推送完成后生成汇总报告包括成功数、失败数、退订过滤数、频控过滤数。8.3 AB 实验与效果分析“多推”的前提是“推得有效”而验证有效性必须靠 AB 实验。最简单的做法推送任务表增加experiment_group字段。把同一批用户随机分成 A、B 两组A 组推优惠文案B 组推稀缺文案最终通过回传数据统计点击率、转化率。SELECT experiment_group, COUNT(DISTINCT push_user_id) AS push_users, COUNT(DISTINCT click_user_id) AS click_users, COUNT(DISTINCT click_user_id) / COUNT(DISTINCT push_user_id) AS ctr FROM push_effect_daily WHERE etl_date 2024-06-20 GROUP BY experiment_group;9. 资源占用与性能观察这套架构的资源占用主要集中在四个环节。9.1 离线计算资源Spark 圈选和画像加工属于离线任务建议提交到 YARN 队列设置队列上限。毕业设计单机环境需要注意内存限制典型参数如下spark-submit \ --master local[4] \ --driver-memory 4g \ --executor-memory 4g \ --conf spark.default.parallelism12 \ job/crowd_select.py如果机器只有 8G 内存驱动和执行器各 4G 可能会直接把机器卡死。更稳妥的做法是先给驱动 2G、执行器 2G跑通了再调大。9.2 Kafka 监控指标实时链路里最容易出现瓶颈的是 Kafka。重点观察两个指标BytesInPerSecKafka 写入速率BytesOutPerSecKafka 消费速率。消费速率长期低于生产速率说明消费者处理不过来可能需要增加分区数或提高消费者并发度。9.3 Redis 频控性能频控判断是高 QPS 场景。很多推送请求会在 Redis 上打点建议使用 pipeline 批量读取用户频控状态避免频繁网络请求。pipe redis_client.pipeline() for user_id in user_list: pipe.get(fpush_count:{user_id}:{today}) pipe.get(flast_push_ts:{user_id}) counts pipe.execute()9.4 数据库瓶颈推送任务表不要在推送高峰期做复杂 SQL 查询。批量任务执行模块应该定时把任务进度写入 Redis 缓存前端页面展示时读取缓存查询明细再从数据库取。否则百万级推送任务表上的COUNT查询会拖垮数据库。10. 常见问题与排查方法问题现象可能原因排查方式解决方案推送任务一直 pending调度器未启动或人群文件未生成查看调度日志和任务表 create_time手动触发调度器确认人群文件就绪用户收到重复推送上游请求未传 request_id查推送日志中的重复 task_id增加 request_id 幂等判断一批用户大量推送失败第三方渠道限流或 Token 过期查看渠道返回的错误码分批推送失败任务自动退避重试频控不生效Redis key 未设置过期时间查看 Redis 中的 key TTL补全过期时间召回结果都是热门内容规则召回未配置个性化特征检查画像表更新时间完善用户最近行为特征Spark 任务内存溢出executor 内存配置过大或过小查看 YARN 日志 OOM 记录调整 executor 内存和并行度效果报表点击率为 0埋点回传数据未关联 task_id查看 Kafka back_topic 日志在推送内容埋点中带上 task_id用户退订后仍收到推送退订表未同步到推送过滤逻辑检查退订表数据更新推送前强制查询退订状态最容易出错的是“退订用户仍然收到推送”这个问题的原因通常是退订记录存在 MySQL但推送服务读取的是 Redis 缓存Redis 缓存没有同步退订状态。退订更新时一定要同步删除或更新 Redis 中的可推送标记。11. 最佳实践与使用建议第一先做最小可用链路。不要一上来就上 Flink、HBase、ClickHouse。用 MySQL Hive Spark Kafka 先把“埋点 → 画像 → 召回 → 推送 → 回传”跑通再逐步加复杂组件。第二把“退订优先”写进代码逻辑。推送服务流程中退订检查永远放在第一步。任何推送策略、任何运营需求都不能越过退订判断。第三批量任务一定要做日志。每条推送任务至少记录用户 ID、任务 ID、推送时间、推送结果、失败原因。这些日志既是排查问题的依据也是后续效果分析的数据来源。第四模板内容管理与测试。推送内容建议走模板管理避免在代码里拼接消息文案。每次模板变更前先用少量白名单用户测试确认文案和跳转链接无误再批量发送。第五设置“灰度推送”机制。新配置上线时先选择 5% 用户灰度推送观察点击率和投诉率再逐步放大比例。第六涉及用户手机号、设备信息、行为轨迹时必须严格遵守个人信息保护相关法律要求。数据在加工过程中要脱敏内部操作要有权限审计用户授权的获取和撤销记录必须完整留存。12. 总结与下一步“大数据这么会推那就多推”听起来像是运营口号落到工程上其实是三件事推荐要准、推送要稳、批量要大。这套架构最值得先验证的环节是“人群圈选 推送频控”。这两个模块做扎实了系统就具备上线的基本条件。最容易踩的坑是内容推荐做得不错但推送环节没有做好退订过滤和频控设计结果高曝光带来的却是高退订。如果继续扩展可以往三个方向走引入更丰富的画像特征比如用户生命周期价值预测、流失概率模型在推荐链路中引入深度排序模型把推送渠道扩展为多通道协作比如 App 推送、短信、站内信、企微统一在触达层做优先级编排。先把第一版跑起来再谈优化。数据回流形成的闭环才是这套系统真正的价值所在。