简介:这是一份基于Python实现的用户画像生成系统源码,适合需要将原始用户数据转化为可落地画像模型的数据分析师、Python开发者及产品运营人员。代码按照真实业务链路组织,覆盖数据清洗与预处理、特征工程、用户行为分析、聚类分群、特征权重计算、画像构建与可视化展示,并给出Flask/Django等Web服务集成思路,能够帮助读者快速上手个性化推荐与精细化运营场景。压缩包共149个文件,整体大小2.45MB,以69个Python脚本为主干,同时包含53个pyc编译文件、5个HTML页面、4个CSS样式表、JS与source map前端资源、2个CSV示例数据集以及图片、字体等辅助文件,便于同时学习后端逻辑与前端展示。目前已有352人学习下载,适合具备一定Python和数据分析基础的人群作为系统实践参考,也可直接作为课程设计或内部推荐系统的原型代码。
1. 用户画像生成系统:靠 Python 把它从业务口号变成能跑的代码
运营提需求永远是一句话:“找出过去 30 天活跃但没下单的高价值用户,给她们推优惠券。” 你打开数据库,发现只有订单表、登录日志和商品浏览记录。所谓“高价值”“活跃”“未转化”,在 Python 的世界里就是一张标签宽表、一堆加工脚本和一个查询接口。用户画像生成系统,本质就是把业务描述翻译成可计算的标签逻辑,再把这些标签组织成可查询、可更新、可追溯的数据资产。它解决的从来不是“要不要做画像”,而是“标签怎么算、存在哪、多久更新一次、查出来准不准”。这事适合两类人:一类是后端或数据开发,要为公司搭一套能用的画像服务;另一类是独立开发者,手上有业务数据,想靠 Python 快速验证画像能带来什么价值。但先泼一盆冷水:用户画像系统的难点根本不在机器学习,而在数据清洗、标签口径和更新策略。博主见过太多项目死在建模之前——标签算出来了,运营不敢用,因为口径对不上。这篇笔记按“标签设计 → 加工脚本 → 存储调度 → 踩坑排查 → 落地验证”的顺序,把一套可复现的 Python 用户画像生成系统完整拆给你。
2. 先把标签体系设计明白:没有业务语义的画像只是一张没人看的宽表
2.1 标签三层结构:事实标签、规则标签、算法标签各管一段
用户画像系统最常见的翻车方式,是一上来就写代码,把几十个字段堆进一张表,然后发现没人知道“value=3”到底代表什么。我一般会先在文档里把标签分成三层,每一层的计算方式和维护责任都不同。
第一层是事实标签,直接从数据表里取,不加加工。例如注册时间、性别、城市、设备型号。这类标签准确率是 100%,因为它是真实记录的映射。但注意,事实标签也可能过期——用户搬家了,城市字段没更新,这时候“事实”就成了“历史事实”。所以事实标签也要带时间戳。
第二层是规则标签,靠明确的业务规则计算。例如 RFM 分层(最近消费时间、消费频率、消费金额)、活跃度分桶、流失预警。这一层是画像系统的主力,也是运营最买账的一层,因为它能说出“为什么给这个人打这个标签”。规则标签的关键不是规则本身,而是阈值设在哪。后面第 3 章会专门讲用分位数而不是拍脑袋定阈值。
第三层是算法标签,靠机器学习或统计模型产出。例如“高潜付费用户”“母婴兴趣人群”“价格敏感型用户”。这一层听起来高级,但落地时最容易变成黑匣子——算法给一个用户打了“高价值”标签,运营一问为什么,你说不出所以然。我的建议是算法标签只做辅助,规则标签做主力,至少在 v1 版本必须这样。系统先跑起来,运营用出信任感,再逐步加模型。
2.2 三张核心表:用户、行为、标签字典怎么用 Python 建模
在设计标签体系之前,先把底层数据模型理顺。一个用户画像系统,在中小数据规模下(百万级用户,千万级行为),三张表就够了:用户维度表、行为事实表、标签结果表。用户维度表存静态属性,行为事实表存每次点击、下单、收藏事件,标签结果表是最终产物——一个用户一行,一行里几十个标签字段。
下面用 SQL 建表,因为标签加工最终大概率要落到查询引擎上,直接写 SQL 建模最贴近生产环境:
-- 用户维度表:静态属性,缓慢变化维 CREATE TABLE dim_user ( user_id STRING COMMENT '全局唯一用户ID,如设备ID或手机号加密值', register_time TIMESTAMP COMMENT '注册时间', gender STRING COMMENT '性别:male/female/unknown', city STRING COMMENT '城市编码,建议用行政区划码', device_brand STRING COMMENT '设备品牌', last_update_time TIMESTAMP COMMENT '该行ETL更新时间' ) PARTITIONED BY (dt STRING COMMENT '数据分区,按天'); -- 行为事实表:每次行为一行,画像的数据源头 CREATE TABLE dwd_user_event ( user_id STRING, event_time TIMESTAMP, event_type STRING COMMENT 'click/view/add_cart/order/collect', item_id STRING COMMENT '商品或内容ID', item_category STRING COMMENT '商品类目ID', session_id STRING COMMENT '会话ID,用于去重和路径分析', extra MAP<STRING,STRING> COMMENT '扩展字段,如停留时长、页面来源' ) PARTITIONED BY (dt STRING); -- 标签结果表:画像系统的最终产物,一个用户一行 CREATE TABLE dws_user_profile ( user_id STRING, label_json STRING COMMENT '标签集合,JSON格式,方便动态扩展', rfm_level STRING COMMENT '高价值/中价值/低价值', active_level STRING COMMENT '高活跃/中活跃/低活跃/沉默', interest_tags ARRAY<STRING> COMMENT '兴趣标签,如["数码","运动"]', etl_time TIMESTAMP COMMENT '本次生成时间' ) PARTITIONED BY (dt STRING);代码背后的逻辑,说三个关键点:第一,用户表不要存所有历史修改,只需要最新快照,历史版本靠“缓慢变化维”另开表,v1 用不着;第二,行为表的 extra 字段用 MAP 类型,因为埋点字段经常变,固定字段一旦加列就要改表结构,MAP 能扛住业务迭代;第三,标签结果表把标签存成 JSON 而不是几十个独立字段,因为标签体系一定会在三个月后加新标签,JSON 可以在不改表结构的前提下动态扩展,代价是查询时不能直接索引某个标签,但配合 Redis 或倒排索引可以解决。
建完表,还要一张标签字典表,这是画像系统最容易漏掉但最救命的一张表。它存的是元数据:标签名、标签值含义、计算逻辑、数据来源、负责人、更新频率。没有字典表的画像系统,三个月后就是一本天书——你根本不知道“active_level=2”是什么意思。用 Python 维护一张 Markdown 或 Excel 字典即可,但必须挂在代码仓库里,随版本迭代。
2.3 标签体系设计检查清单:动手写代码前先回答三个问题
每个标签在写计算逻辑之前,强迫自己回答三个问题:这个标签谁来用、用到什么场景、标签错了会怎样?第一个问题决定标签存储格式——运营要的人群圈选标签,要可检索;算法要的用户特征,要向量化。第二个问题决定更新频率——活动营销标签要小时级更新,用户价值分层可以天级更新。第三个问题决定准确率要求——给用户打“高价值”标错,最多浪费一张优惠券;给用户打“疑似欺诈”,标错就要出客诉。
我在 v1 阶段只保留高频使用的 15 到 20 个标签,不追求大而全。画像系统最重要的不是标签数量,而是每个标签都能被业务说清楚。先跑通核心链路,再慢慢加。把标签体系定下来,接下来进入最关键的环节——用 Python 把原始轨迹变成可靠标签。
3. 用 Python 把行为数据加工成标签:ETL 脚本与特征计算的完整实现
3.1 埋点日志清洗:去重、过滤爬虫、时间对齐三板斧
画像系统的数据源通常是埋点日志。埋点日志的质量比想象中差得多,不处理直接加工,标签准确率会低到你怀疑人生。清洗部分我一般做三件事:去重、过滤无效用户、时间对齐。
去重不是简单地删掉一模一样的行。用户在前端快速点击,同一个事件可能被上报两次;App 离线时用户产生的事件,恢复网络后可能重传。所以去重要按“用户 ID + 会话 ID + 事件类型 + 事件时间戳”四元组来判重。时间对齐更隐蔽——客户端时间和服务端时间可能差了几分钟,如果直接用客户端时间做“24小时活跃”的判断,边缘用户会被错分。
下面给出一段可运行的清洗代码,用 Pandas 处理一批日志文件:
# -*- coding: utf-8 -*- """埋点日志清洗:幂等,可重复执行""" import pandas as pd def clean_event_log(df: pd.DataFrame) -> pd.DataFrame: """ 输入字段必须包含: user_id, session_id, event_type, event_time, client_time """ # 1. 时间对齐:以服务端接收时间为准, # 客户端与服务端时间差超过 5 分钟则以后者为准 time_diff = (df["event_time"] - df["client_time"]).dt.total_seconds() df.loc[time_diff.abs() > 300, "event_time"] = df.loc[time_diff.abs() > 300, "client_time"] # 2. 去重:四元组完全一致视为同一条事件 df = df.drop_duplicates( subset=["user_id", "session_id", "event_type", "event_time"] ) # 3. 过滤异常 user_id:长度小于 5 或者是 NaN 的直接丢掉 df = df[df["user_id"].notna() & (df["user_id"].astype(str).str.len() >= 5)] # 4. 过滤爬虫/测试流量:UA 特征或 event_type=click 却每秒超过 20 次 click_count = df[df["event_type"] == "click"].groupby("user_id").size() spider_users = click_count[click_count > 500].index # 单日 500 次点击以上先标记 df = df[~df["user_id"].isin(spider_users)] return df这段代码里的几个参数值得说明。时间对齐的 300 秒是经验值——超过 5 分钟的时间偏差,大概率是用户改了系统时间,而不是正常网络延迟,正常延迟在 10 秒以内即可。过滤爬虫的 500 次点击阈值要根据业务调,内容社区用户一天点几百次很正常,工具类产品用户一天点五十次就算重度。先标出来人工复核,不要直接删,因为有些渠道会刷量。Pandas 在百万行级以内性能没问题,到了千万行级就要换 PySpark 或 DuckDB,后面讲。
清洗完的数据,才能进入标签加工环节。
3.2 规则标签计算:RFM 分层不是拍脑袋,用分位数找阈值
RFM 是用户画像里最经典、运营最容易认可的规则标签。但网上绝大多数 RFM 教程都写了错的做法:R 小于 7 天算高价值、F 大于 10 次算高价值、M 大于 500 元算高价值——这些阈值是拍脑袋拍的。不同业务的数据分布天差地别,生鲜电商的用户每月买 20 次是常态,二手车平台用户三年买一次也正常。RFM 阈值必须从数据里算出来。
正确的做法是分位数切分。R(最近消费时间)取倒数值或者直接用负号,这样三个维度都是“越大越好”。然后用分位数把每个维度切成两段或三段。下面的代码实现了一个完整 RFM 计算流程:
# -*- coding: utf-8 -*- """RFM 分层计算脚本,基于订单表""" import pandas as pd import numpy as np def compute_rfm(orders: pd.DataFrame, reference_date: str = None) -> pd.DataFrame: """ orders 必须字段: user_id, order_time, order_amount reference_date: 计算基准日,不传则取订单表最大日期 """ if reference_date is None: reference_date = orders["order_time"].max() reference_date = pd.Timestamp(reference_date) # 按用户聚合:R 是最近一次下单距今天数,F 是下单次数,M 是累计金额 rfm = orders.groupby("user_id").agg( R=("order_time", lambda x: max(0, (reference_date - x.max()).days)), F=("order_time", "count"), M=("order_amount", "sum") ).reset_index() # 用分位数切分:Q1 以下为低分段(0),以上为高分段(1) # 这里用了 50% 分位做二分类,可以按业务调成 33%/67% rfm["R_score"] = (rfm["R"] <= rfm["R"].quantile(0.5)).astype(int) rfm["F_score"] = (rfm["F"] > rfm["F"].quantile(0.5)).astype(int) rfm["M_score"] = (rfm["M"] > rfm["M"].quantile(0.5)).astype(int) # 三段式:111 高价值 / 011 发展 / 100 流失预警 def rfm_level(row): s = f"{row['R_score']}{row['F_score']}{row['M_score']}" mapping = { "111": "高价值", "011": "潜力用户", "101": "新客高消", "001": "重要挽留", "110": "活跃低消", "010": "一般用户", "100": "流失预警", "000": "低价值", } return mapping.get(s, "未知") rfm["rfm_level"] = rfm.apply(rfm_level, axis=1) return rfm[["user_id", "R", "F", "M", "R_score", "F_score", "M_score", "rfm_level"]]这段代码最关键的一行是quantile(0.5)。它把用户从“跟绝对标准比”变成“跟大盘比”,这是画像系统能落地的核心。但分位数切分有个坑:如果某个维度的数据分布极端,比如 80% 的用户消费金额是 0,那么 M 的分位数切分会把所有非零用户都归为高分,区分度反而变差。碰到这种情况,我一般先过滤掉无消费用户,再在消费用户群体内部做分位。
R 的取值方向要注意——R 是最近一次消费距今的天数,所以“R 越小越好”。代码里用了<=判断算 1 分,含义是“R 小于大盘中位数,说明最近刚消费过,得分高”。这段逻辑写得隐晦,建议在生产代码里加注释。
3.3 活跃度分桶:把 DAU 定义具象成可配置参数
活跃度标签是画像系统最高频的字段,但“活跃”的定义在不同业务里完全不同——内容产品用“打开次数”定义,电商用“浏览商品数”,工具产品用“核心功能使用次数”。我一般把活跃度做成参数可配置的分桶逻辑,而不是硬编码。
def compute_active_level(user_events: pd.DataFrame, config: dict) -> pd.DataFrame: """ user_events: user_id, ds(活跃日期), event_type config 示例: { "high": {"min_days": 7, "min_events": 30}, # 7天内活跃>=7天且事件>=30次 "mid": {"min_days": 3, "min_events": 10}, "low": {"min_days": 1, "min_events": 1}, } """ # 最近 7 天活跃天数和总事件数 recent = user_events[user_events["ds"] >= (user_events["ds"].max() - pd.Timedelta(days=6))] stats = recent.groupby("user_id").agg( active_days=("ds", "nunique"), total_events=("event_type", "size") ).reset_index() def _level(row): if row["active_days"] >= config["high"]["min_days"] and row["total_events"] >= config["high"]["min_events"]: return "高活跃" if row["active_days"] >= config["mid"]["min_days"] and row["total_events"] >= config["mid"]["min_events"]: return "中活跃" if row["active_days"] >= config["low"]["min_days"]: return "低活跃" return "沉默" stats["active_level"] = stats.apply(_level, axis=1) return stats[["user_id", "active_days", "total_events", "active_level"]]注意代码里的“近 7 天”是相对计算日往前推 6 天,这是写增量任务的常见坑——很多人写now()去取日期,结果重跑历史分区的数据时,取到的是跑批当天的日期,导致历史分区的活跃度全错。正确的做法是计算日从上文传入,或者从分区字段里取。
3.4 算法标签:用 TF-IDF 从行为序列提取用户兴趣词
规则标签做完,画像系统已经能回答“用户值不值得运营”的问题,但还回答不了“用户对什么感兴趣”。这时需要一个轻量级的算法标签——从用户浏览过的商品标题或内容标题里提取兴趣关键词。这里不用上深度学习,TF-IDF 或词频统计足够支撑 v1,跑得飞快且可解释。
# -*- coding: utf-8 -*- """基于 TF-IDF 的用户兴趣标签提取""" from sklearn.feature_extraction.text import TfidfVectorizer def extract_interest_tags(item_titles: dict) -> dict: """ item_titles: {user_id: ["商品标题1", "商品标题2", ...]} """ # 每个用户看过的所有标题拼接成一篇“文档” user_ids = list(item_titles.keys()) documents = [" ".join(item_titles[uid]) for uid in user_ids] if not documents: return {} # 停用词按业务自定义,电商场景要加“包邮”“正品”等噪声词 vec = TfidfVectorizer( token_pattern=r"(?u)\b[\w\u4e00-\u9fa5]{2,}\b", # 匹配中文词和英文词 max_features=2000, min_df=3, # 低于 3 个用户出现的词直接丢弃 max_df=0.5, # 超过 50% 用户出现的词是噪声,丢弃 stop_words=None # 业务停用词表自己传 ) tfidf = vec.fit_transform(documents) # 每行取 top3 词作为兴趣标签 result = {} feature_names = vec.get_feature_names_out() for idx, uid in enumerate(user_ids): row = tfidf[idx].toarray()[0] top_indices = row.argsort()[-3:][::-1] tags = [feature_names[i] for i in top_indices if row[i] > 0.01] result[uid] = tags return resultTfidfVectorizer 的max_df=0.5是值得注意的参数——如果一个词出现在超过一半用户的历史里,那这个词就是“性价”“质量”这类没有区分度的泛词。同理,min_df=3过滤了低频噪声词,比如商品标题里的错别字或超长尾型号数字。max_features=2000是内存保护。
这一段代码的逻辑可以变体为:把商品标题换成文章标题、短视频 tag、搜索关键词,就能适配不同业务的兴趣标签。但本质上,它是“用户看过的东西 → 内容关键词”的共现统计,属于最浅层的兴趣模型。如果用户行为序列足够长(人均上百次浏览),下一步可以换成 Word2Vec 或 Sentence-BERT 做语义级标签。
3.5 标签合并与宽表落库:幂等写入是底线
以上三段代码分别产出了 RFM 标签、活跃度标签和兴趣标签,最后一步是把它们合并成画像系统的核心产物——用户宽表。这一步有两个硬性要求:一是幂等,同一批数据无论跑多少次,结果都完全一致;二是可追溯,标签从哪个计算脚本哪个版本产出的,必须能查。
# 合并完全部标签到 user_profile DataFrame # profile_df 至少有 user_id, label_json, etl_time profile_df["label_json"] = profile_df.apply( lambda row: json.dumps({ "rfm_level": row["rfm_level"], "active_level": row["active_level"], "interest_tags": row["interest_tags"], }, ensure_ascii=False), axis=1 ) # 写入 MySQL(或 Hive / ClickHouse),用 ON DUPLICATE KEY UPDATE 保证幂等 from sqlalchemy import create_engine, text engine = create_engine("mysql+pymysql://user:pass@host:3306/profile_db") with engine.begin() as conn: for _, row in profile_df.iterrows(): conn.execute(text(""" INSERT INTO dws_user_profile (user_id, label_json, etl_time, dt) VALUES (:user_id, :label_json, NOW(), :dt) ON DUPLICATE KEY UPDATE label_json=VALUES(label_json), etl_time=NOW() """), {"user_id": row["user_id"], "label_json": row["label_json"], "dt": biz_date})注意写入用的ON DUPLICATE KEY UPDATE——同一天的跑批如果重复执行,不会产生重复行,而是覆盖更新。这就是幂等。但逐行iterrows在数据量大时性能很差,实际生产一般用pd.to_sql先写入临时表,再INSERT INTO ... SELECT或者直接REPLACE INTO。单行插入这段只是把逻辑讲清楚,真正做的时候换成批量方式。
标签加工阶段跑通,画像系统有了“内容”。下一章解决“怎么让它在业务上稳定跑起来”——存储选型和更新策略。
4. 画像存储与服务化:宽表查询、Redis 缓存和定时更新策略
4.1 存储方案选型:MySQL 宽表、Redis 哈希、ClickHouse 的适用边界
用户画像数据怎么存,取决于你的查询场景。最典型的两类查询:一是“给我这 100 万用户的详细标签”,二是“圈出所有高价值高活跃且兴趣是数码的用户”。第一类查宽表就行,第二类需要标签索引。
下面这个表格是我在不同项目里的选型参考:
| 存储方案 | 适用场景 | 单用户标签数 | 查询 QPS | 运维成本 | 典型数据量 |
|---|---|---|---|---|---|
| MySQL 单表/分片 | 千万级用户以下,按 user_id 点查 | 50 以内 | 数百 | 最低 | 千万行以内 |
| Redis 哈希 | 在线实时查询,延迟要求 <10ms | 任意,JSON 序列化 | 数万 | 低 | 百万到千万用户 |
| ClickHouse | 人群圈选、标签组合筛选、Ad-hoc 分析 | 50 以上,列式存储 | 数百 | 中 | 亿级以上 |
| HBase | 超大规模用户,稀疏标签矩阵 | 任意 | 数千 | 高 | 十亿级 |
v1 阶段我默认推荐 MySQL 宽表 + Redis 缓存。理由很简单:MySQL 做持久化和回溯重算足够,Redis 做线上查询的加速层,两个组件你的技术栈里大概率已经有了。ClickHouse 是人群圈选的最佳选择,但引入一个新的分布式组件意味着运维成本上升。公众号文章里那些大厂画像体系动辄 HBase + Spark 的架构,你只有几千万数据,真没必要。
4.2 画像数据如何进 Redis:按 user_id 做哈希,还是按标签做倒排
Redis 在画像系统里有两种用法,对应的数据结构完全不同。
第一种是按用户维度缓存。key = profile:{user_id},value 是用户全量标签的 JSON 串。查询画像时直接 GET 就返回,延迟在毫秒级。这种做法的优点是实现简单,适合“单用户画像展示”的场景。缺点是它不能回答“有哪些用户是高价值”——要回答这个问题,你得扫全量 key,在 Redis 里这就是灾难。
第二种是按标签维度做倒排索引。key = tag:rfm_level:高价值,value 是 Set,里面存这个标签覆盖的所有 user_id。这样圈人时直接SINTER或SUNION多个 Set,一次交并集运算就出人群包。代价是更新麻烦——标签变了,你不但要更新用户的 profile key,还要维护倒排索引的一致性。
我的做法是两层同时用:MySQL 宽表算新鲜数据,Redis 用 pipeline 批量写 profile key 支撑在线查询,再用一个异步任务把高频圈人标签同步成倒排 Set。更新顺序要注意,必须先更新 MySQL,等数据落库后再刷新 Redis,这样万一 Redis 挂了,MySQL 还能兜底恢复。
下面给一段把宽表批量灌入 Redis 的核心代码:
import redis r = redis.Redis(host="localhost", port=6379, db=1) def sync_profile_to_redis(profile_df) -> None: """将画像宽表同步到 Redis Hash,使用 pipeline 提升写入性能""" pipe = r.pipeline(transaction=False) for _, row in profile_df.iterrows(): key = f"profile:{row['user_id']}" # HSET 的好处是后续标签字段变化,只更新对应 field,不需要覆盖整个 JSON pipe.hset(key, mapping={ "rfm_level": row["rfm_level"], "active_level": row["active_level"], "interest_tags": row["interest_tags"], "etl_time": str(row["etl_time"]), }) pipe.expire(key, 3600 * 24 * 3) # 设置过期时间,防止用户长期不更新占内存 pipe.execute()注意这里用的是HSET而不是字符串整体写入。原因是一个用户的标签经常只变其中一个字段(比如活跃度从“高”变“中”),用 Hash 结构可以只更新一个 field,减少不必要的网络开销。expire设了 3 天,意思是用户 3 天没有新数据,缓存自动失效,下次查询时回源 MySQL 重新加载。这是 Redis 缓存最常见的坑——不设过期时间的缓存,数据永远不同步。
4.3 定时更新策略:全量重算和增量更新怎么搭配
画像标签有各自不同的更新频率。用户性别、注册城市基本不变,月级重算一次足够;RFM 里的“最近消费时间”每天都会变,要天级重算;活动营销用的“今日活跃用户”时效性最强,需要小时级甚至准实时。
一个务实的分层策略是这样:
| 标签类型 | 更新频率 | 计算方式 | 说明 |
|---|---|---|---|
| 事实标签(性别/城市) | 每周 | 全量核对源表 | 源表变了才更新,源表不变跳过 |
| RFM 分层 | 每天凌晨 | 增量拉当天订单 + 全量聚合 | R 必须全量算,因为时间维度变了 |
| 活跃度 | 每小时 | 增量 append | 只统计最近 7 天窗口 |
| 兴趣标签 | 每天 | 增量 + 全量合并 | 新行为只影响近期的兴趣权重 |
| 临时活动标签 | 按需 | 实时计算 | 活动结束即销毁 |
实际操作用 APScheduler 写定时任务:
# 定时任务配置示例 from apscheduler.schedulers.blocking import BlockingScheduler from apscheduler.triggers.cron import CronTrigger scheduler = BlockingScheduler() # 每天凌晨 2 点重算 RFM scheduler.add_job( compute_rfm_job, CronTrigger(hour=2, minute=0), id="rfm_daily", replace_existing=True, max_instances=1 # 关键参数:防止上一次没跑完,下一次又启动 ) # 每小时第 15 分钟更新活跃度 scheduler.add_job( compute_active_job, CronTrigger(minute=15), id="active_hourly", replace_existing=True, max_instances=1 ) scheduler.start()max_instances=1是定时任务的血泪教训——不设置的话,一旦任务执行时间超过了调度间隔,同一时刻会有两个实例在跑,轻则重复计算浪费资源,重则互相抢锁把目标表写坏。
定时任务除了调度器,还要有两个保障机制:一个是血缘记录,在每次计算后往etl_job_log表里写一行日志,记录“哪个脚本、哪个版本、处理了多少用户、耗时多久、产出多少标签”;另一个是失败重跑,失败的任务要能手动单独重跑,而不是重新触发全流程。这两个机制在出现数据事故时是后悔药。
5. 用户画像的常见坑与排查:标签覆盖率、数据漂移和内存翻车
5.1 标签覆盖率骤降:从一个 90% 掉到 40% 的排查过程
现象:某天运营反馈“高价值用户”人数突然少了 60%,之前几周都正常。检查调用方代码、查询口径、Redis key 都没有变化,最后用 SQL 查了标签结果表,发现label_json字段大量为NULL。
原因:前一天上线新版埋点 SDK,order_time字段的格式从字符串"2024-01-15 12:00:00"改成了时间戳1705305600。RFM 计算脚本里pd.Timestamp(reference_date) - x.max()变成负值,但代码里没有拦截,R取max(0, ...)之后全部变成 0,随后R_score大量为 0,高价值标签自然崩了。
解决:在计算脚本入口加字段格式校验——assert或者抛异常,宁可让任务失败也不要产出脏数据。另外,给所有标签加工脚本加了“覆盖率监控”,每次跑批后自动统计label_json IS NOT NULL的用户比例,低于阈值直接告警,而不是等到运营来反馈。
5.2 用户画像的标签漂移:昨天还是“高活跃”,今天变成“沉默”
现象:同一批用户的活跃度标签在前后两天没有任何真实行为变化的情况下,出现了 9% 的用户从“高活跃”跌到“沉默”。查活跃底表,没有数据丢失;查刷新接口,Redis 缓存时间正常。
原因:“近 7 天活跃窗口”是按绝对日期算的。哪天是这一周的新起点,窗口边缘用户的活跃天数会被“挤掉”一天,如果他们的活跃天数恰好等于阈值下限,就掉桶了。更隐蔽的是,跑批时用了datetime.now()取当前时间,一旦任务延迟到凌晨 1 点跑,和历史分区数据的“当天”语义不一致。
解决:所有标签计算的基准日期都用“调度日期”显式传入,不允许从系统时间获取。活跃度窗口不是“最近 7 天”,而是“计算日往前推 6 天”,这样窗口随调度日期平滑移动,不再有跳变。
5.3 内存翻车:Pandas 处理千万级行为日志直接 OOM
现象:本地测试数据量 500 万行没问题,上线后数据量到了 3000 万行,ETL 脚本跑了几分钟被系统 kill。回头看代码,全程 Pandas DataFrame,光读取日志就占掉 24GB 内存。
原因:Pandas 是内存计算框架,数据量一上来就撑不住。另外代码里有groupby之后apply自定义 lambda,每个分组都要复制一份数据,内存翻倍。
解决:千万行级日志换用 DuckDB 或 PySpark。PySpark 学习成本高,DuckDB 可以保持 SQL/Pandas 风格无缝切换。如果必须在内存里算,就分块读取,每 200 万行处理一次,合并结果,并且groupby里的 apply 全部改成向量化写法。博主踩过一次iterrows写数据库的坑,300 万行用户逐行写 MySQL 跑了 40 分钟,后来换成 DataFrame 批量写临时表再 SELECT 的方式,只用了 2 分钟。
5.4 用户 ID 口径混乱:画像标签全串号了
现象:画像表里一个用户行上,既有手机号加密值作为 user_id,又有设备 ID 作为 user_id,两套体系同时灌入,导致同一实体在表里出现两行甚至两行标签互相矛盾。
原因:埋点 SDK 的 user_id 在用户未登录的时候取的是设备 ID,登录后取的是账号 ID。两条事件可能被当成两个不同的人。画像系统没有做 ID 归一。
解决:在加工入口做 ID-Mapping。用一个映射表id_mapping把设备 ID、账号 ID、手机号加密值全部映射到统一的global_user_id。业务主键在映射表里维护,多对一关系。做一次全量映射更新,后续新进来的行为日志都先查映射表再进画像流程。
5.5 特征穿越:拿未来数据给过去打标签
现象:某次活动复盘,发现画像系统给“双 11 之前一周”的用户打上了“已购买用户”标签,运营看着这批用户感觉不对劲,一查才发现逻辑漏洞。
原因:标签计算任务拉了全量行为表,没有按行为事件时间过滤。用户在 11 月 12 日下单,而任务重算的是 11 月 5 日那天的分区,数据混在一起,标签产生了穿越。
解决:行为表中所有事件都有event_time,标签计算必须按照事件时间过滤,而不是跑批日期过滤。换句话说,用户画像的每一个标签,都只能基于“标签生效时刻”之前已经发生的数据。写计算逻辑的时候,把“最晚可用事件时间”作为一个参数传进去,跑历史分区时就传历史分区日期。
6. 画像真正用起来:人群圈选 API、画像报告和效果验证
6.1 搭建一个最简人群圈选 API:用 FastAPI 暴露标签查询
画像系统最终要变成业务系统能调用的服务。最常用的接口就是两个:按 user_id 查询画像、按标签组合圈人群包。查画像直接走 Redis;圈人群走 MySQL 的 JSON 字段查询,数据量大时再用 ClickHouse 优化。下面是 v1 版人群圈选接口:
from fastapi import FastAPI, Query import json from sqlalchemy import text app = FastAPI(title="User Profile Service") @app.get("/api/v1/profiles/query") def query_profiles( rfm_level: str = Query(None, description="高价值/潜力用户/流失预警等"), active_level: str = Query(None, description="高活跃/中活跃/低活跃/沉默"), interest: str = Query(None, description="兴趣标签,单值"), limit: int = Query(100, le=1000), ): """ 标签组合查询,返回符合条件的 user_id 列表。 注意:这里直接查 MySQL 的 label_json 字段, 数据量大时要用倒排索引或 ClickHouse 替换,否则扛不住。 """ conditions = [] params = {"limit": limit} if rfm_level: conditions.append("JSON_EXTRACT(label_json, '$.rfm_level') = :rfm_level") params["rfm_level"] = rfm_level if active_level: conditions.append("JSON_EXTRACT(label_json, '$.active_level') = :active_level") params["active_level"] = active_level if interest: conditions.append("JSON_CONTAINS(JSON_EXTRACT(label_json, '$.interest_tags'), JSON_QUOTE(:interest))") params["interest"] = interest if not conditions: return {"code": 400, "msg": "至少指定一个标签条件"} sql = f"SELECT user_id FROM dws_user_profile WHERE dt=:dt AND {' AND '.join(conditions)} LIMIT :limit" params["dt"] = params.get("dt", latest_partition) with engine.connect() as conn: rows = conn.execute(text(sql), params).fetchall() return {"code": 0, "data": [row[0] for row in rows]}这个接口上线时要注意两个问题。JSON_EXTRACT 在百万级表上还能跑,到了千万级就要加 MySQL 生成的虚拟列 + 索引,或者干脆导出倒排索引到 Redis,毕竟全表扫 JSON 是数据库的大忌。接口参数要先做白名单校验,否则用户传一个“高价值’ OR 1=1”进来,直接构成 SQL 注入。要用参数化查询,不要拼字符串。
6.2 验证画像的准确率:单标签抽检和 AB 实验结合起来
画像标签算完没人验证,等于白做。验证分两层。
第一层是单标签准确率抽检。比如“高价值用户”标签,抽样 200 个被标记为高价值的用户,人工核对她们的最近订单记录,看是否名实相符。准确率超过 85% 算合格。这里的 85% 是我自己的标准,解释成本低。算法标签要额外看覆盖率——如果“兴趣标签”只能覆盖 30% 的用户,那这个标签对另外 70% 的用户没有意义,要么放弃,要么用默认值兜底。
第二层是业务效果的 AB 实验。画像的价值最终体现在“用了画像 vs 不用画像,业务指标差多少”。比如运营做新客转化活动,对照组用随机用户,实验组用画像圈出来的“高活跃但未下单”用户。观察两组同样的优惠券发放成本和订单转化率。如果实验组转化率没有显著提升,画像标签再精致也是无效资产。这个验证周期一般是 2 到 4 周。
6.3 一个进阶技巧:用规则引擎把标签改动从代码里解放出来
画像标签的改动频率远比你想象的高。运营每个月都要加新标签、调旧阈值,如果每次改标签计算逻辑都要改代码、重新部署,开发会累死,运营也会失去耐心。我后期把规则标签的阈值配置放进了数据库表:
CREATE TABLE tag_rule_config ( tag_name STRING COMMENT '标签名,如 active_level', rule_param STRING COMMENT '参数名,如 min_days', rule_value STRING COMMENT '参数值,如 7', update_time TIMESTAMP );计算脚本每次跑批前先读这个配置表,把参数动态注入计算逻辑。这样运营调整“高活跃”的阈值时,只需要改一行配置,完全不需要动代码。这个模式救了我很多次,强烈建议在 v1 阶段就加上。
最后的最后说个习惯:画像系统的每张表和每个标签,我都会在文档里有一行“期望数据波动范围”。比如 RFM 标签里的“高价值”用户占比,正常区间是 10% 到 20%。每天跑批后看一眼核心标签的分布变化,一旦超出预期范围,立刻停下来查原因,不要等到月底才发现数据已经跑偏两个月。用户画像系统不是搭完就完事的工程,它更像一棵需要持续修剪的树——标签体系会根据业务调整,数据质量会随着埋点升级而变化,存储方案会跟着数据量演进。把这套系统当作一个持续迭代的产品来做,不要让它在角落里自己长草。希望这些思路对你有点帮助,祝你在自己项目里把画像落地成功。
本文还有配套的精品资源,点击获取