欢迎光临
我们一直在努力

RFM 模型工程化落地:用户分层从 Excel 到自动化标签系统

RFM 模型工程化落地:用户分层从 Excel 到自动化标签系统

大家好,我是朱大喜!今天聊一个数据分析师绕不开的经典话题——RFM 模型。别看它简单,真正要在生产环境中做到自动化、可维护、可回溯,里面的坑一个都不少。

一、从 Excel 手搓到自动化需求的诞生

RFM 模型大概是每个数据分析师入行时都会接触的东西。所谓 RFM,就是 Recency(最近一次消费时间)、Frequency(消费频率)、Monetary(消费金额)三个维度。用这三个维度给用户打分,然后做交叉分层,就能把用户分成"重要价值客户"、"重要挽回客户"、"一般维持客户"等不同等级。

刚开始做用户分层的时候,大家可能都经历过这个流程:从数据库导出一张 CSV,扔进 Excel,用 PERCENTILE 函数算分位点,VLOOKUP 分档,再手动标颜色。一周一次勉强能撑住,但要做到每天更新、多维度交叉、自动下发标签,Excel 就完全不够用了。

我们团队接到的需求是:为全平台 2000 万月活用户每天生成 RFM 标签,喂给下游的推荐系统和 CRM 平台。标签要求分层 5 档(R1-R5, F1-F5, M1-M5),每档的切割点基于前 90 天数据的百分位动态计算,还要支持按业务线(电商、内容付费、会员订阅)分别计算。

这个数据量和复杂度,决定了必须走系统化、工程化的路线。

二、RFM 指标计算的 SQL 实现

在数据仓库中,RFM 计算属于典型的窗口分析场景。我们用 Hive SQL 来实现每日增量计算。核心思路是:基于过去 90 天的订单数据,计算每个用户的 R、F、M 原始值,然后映射到 1-5 的标签档位。

— ========== 第一步:计算用户级别的 R、F、M 原始值 ==========
— 基于过去90天的订单数据,按用户聚合计算
WITH user_rfm_raw AS (
SELECT
user_id,
— R值:最近一次消费距今天数(越小越好,代表越活跃)
DATEDIFF(CURRENT_DATE(), MAX(order_date)) AS recency_days,
— F值:近90天消费次数(越大越好)
COUNT(DISTINCT order_id) AS frequency,
— M值:近90天消费总金额(越大越好)
SUM(pay_amount) AS monetary,
— 业务线标识,用于分区计算分位点
business_line
FROM dwd_order_detail
WHERE order_date >= DATE_SUB(CURRENT_DATE(), 90) — 取近90天数据
AND order_status = 'PAID' — 只算已支付订单
AND pay_amount > 0 — 过滤退款等异常数据
GROUP BY user_id, business_line
),

— ========== 第二步:按业务线计算分位点 ==========
— 使用 PERCENTILE 函数计算各指标在 20/40/60/80 分位的值
rfm_quantiles AS (
SELECT
business_line,
— R 值分位点(注意 R 值反向:越小越好,所以高分段对应小值)
PERCENTILE(recency_days, 0.2) AS r_p20,
PERCENTILE(recency_days, 0.4) AS r_p40,
PERCENTILE(recency_days, 0.6) AS r_p60,
PERCENTILE(recency_days, 0.8) AS r_p80,
— F 值分位点
PERCENTILE(frequency, 0.2) AS f_p20,
PERCENTILE(frequency, 0.4) AS f_p40,
PERCENTILE(frequency, 0.6) AS f_p60,
PERCENTILE(frequency, 0.8) AS f_p80,
— M 值分位点
PERCENTILE(monetary, 0.2) AS m_p20,
PERCENTILE(monetary, 0.4) AS m_p40,
PERCENTILE(monetary, 0.6) AS m_p60,
PERCENTILE(monetary, 0.8) AS m_p80
FROM user_rfm_raw
GROUP BY business_line
),

— ========== 第三步:映射 RFM 标签(1-5分) ==========
user_rfm_label AS (
SELECT
a.user_id,
a.business_line,
a.recency_days,
a.frequency,
a.monetary,
— R标签:recency越小越好,所以值越小分越高
CASE
WHEN a.recency_days <= q.r_p20 THEN 5
WHEN a.recency_days <= q.r_p40 THEN 4
WHEN a.recency_days <= q.r_p60 THEN 3
WHEN a.recency_days <= q.r_p80 THEN 2
ELSE 1
END AS r_label,
— F标签:frequency越大越好
CASE
WHEN a.frequency >= q.f_p80 THEN 5
WHEN a.frequency >= q.f_p60 THEN 4
WHEN a.frequency >= q.f_p40 THEN 3
WHEN a.frequency >= q.f_p20 THEN 2
ELSE 1
END AS f_label,
— M标签:monetary越大越好
CASE
WHEN a.monetary >= q.m_p80 THEN 5
WHEN a.monetary >= q.m_p60 THEN 4
WHEN a.monetary >= q.m_p40 THEN 3
WHEN a.monetary >= q.m_p20 THEN 2
ELSE 1
END AS m_label,
— 组合标签:如 "R5_F4_M3"
CONCAT('R',
CASE WHEN a.recency_days <= q.r_p20 THEN '5'
WHEN a.recency_days <= q.r_p40 THEN '4'
WHEN a.recency_days <= q.r_p60 THEN '3'
WHEN a.recency_days <= q.r_p80 THEN '2'
ELSE '1' END,
'_F',
CASE WHEN a.frequency >= q.f_p80 THEN '5'
WHEN a.frequency >= q.f_p60 THEN '4'
WHEN a.frequency >= q.f_p40 THEN '3'
WHEN a.frequency >= q.f_p20 THEN '2'
ELSE '1' END,
'_M',
CASE WHEN a.monetary >= q.m_p80 THEN '5'
WHEN a.monetary >= q.m_p60 THEN '4'
WHEN a.monetary >= q.m_p40 THEN '3'
WHEN a.monetary >= q.m_p20 THEN '2'
ELSE '1' END
) AS rfm_label
FROM user_rfm_raw a
JOIN rfm_quantiles q ON a.business_line = q.business_line
)

— ========== 第四步:映射用户分层 ==========
— 将 RFM 标签组合映射为业务可理解的分层名称
SELECT
user_id,
business_line,
rfm_label,
CASE
WHEN r_label >= 4 AND f_label >= 4 AND m_label >= 4
THEN '重要价值客户'
WHEN r_label >= 4 AND f_label <= 2 AND m_label >= 4
THEN '重要挽回客户'
WHEN r_label <= 2 AND f_label >= 4 AND m_label >= 4
THEN '重要保持客户'
WHEN r_label >= 4 AND f_label >= 4 AND m_label <= 2
THEN '重要发展客户'
WHEN r_label <= 2 AND f_label <= 2 AND m_label <= 2
THEN '流失客户'
WHEN r_label >= 3 AND f_label >= 3 AND m_label >= 3
THEN '一般价值客户'
ELSE '一般维持客户'
END AS user_tier,
CURRENT_TIMESTAMP() AS tag_generate_time
FROM user_rfm_label;

这个 SQL 脚本每天跑一次,处理 2000 万用户的数据大约需要 8 分钟(Hive on Spark,集群 50 个 Executor)。分位点虽然用 PERCENTILE 函数有点暴力,但在百万级以上数据集上,误差通常在 0.5% 以内,业务完全可接受。

三、Python 调度脚本与自动化

SQL 算出了标签,但还需要调度编排、数据校验、异常告警一套流程。我们用 Python + Airflow 搭了一套完整的标签生产 Pipeline。

import pymysql
from datetime import datetime, timedelta
import smtplib
from email.mime.text import MIMEText

# ========== RFM 标签生产 Pipeline ==========
class RFMTagPipeline:
"""RFM 标签自动化生产与校验流程"""

def __init__(self, db_config):
self.db_config = db_config
self.conn = None

def connect_db(self):
"""建立数据库连接"""
self.conn = pymysql.connect(**self.db_config)

def check_data_volume(self, target_date):
"""
数据量校验:对比当日产出与近7日均值
如果偏差超过20%,触发告警
"""
sql = """
SELECT
COUNT(DISTINCT user_id) AS today_cnt,
— 计算近7天(不含今天)的日均用户数
(SELECT AVG(cnt) FROM (
SELECT COUNT(DISTINCT user_id) AS cnt
FROM dwd_user_rfm_label
WHERE dt >= DATE_SUB(%s, 7) AND dt < %s
GROUP BY dt
) t) AS avg_7d_cnt
FROM dwd_user_rfm_label
WHERE dt = %s
"""

with self.conn.cursor() as cursor:
cursor.execute(sql, (target_date, target_date, target_date))
today_cnt, avg_cnt = cursor.fetchone()

# 偏差率计算
deviation = abs(today_cnt – avg_cnt) / avg_cnt if avg_cnt else 0
return today_cnt, deviation

def check_distribution(self, target_date):
"""
分层分布校验:各分层用户占比波动不超过5%
防止因上游数据问题导致标签大面积错乱
"""
sql = """
SELECT
user_tier,
COUNT(*) AS user_cnt,
COUNT(*) * 100.0 / SUM(COUNT(*)) OVER() AS pct
FROM dwd_user_rfm_label
WHERE dt = %s
GROUP BY user_tier
ORDER BY user_cnt DESC
"""

with self.conn.cursor() as cursor:
cursor.execute(sql, (target_date,))
return cursor.fetchall()

def send_alert(self, subject, content):
"""
告警通知:企业微信机器人 / 邮件
"""
# 使用企业微信 Webhook 发送告警
# 这里简化展示为邮件方式
msg = MIMEText(content, 'plain', 'utf-8')
msg['Subject'] = f'[RFM标签告警] {subject}'
msg['From'] = 'data-platform@company.com'
msg['To'] = 'data-team@company.com'
# smtp.send_message(msg) # 实际使用需配置 SMTP

def run(self, target_date):
"""主流程:执行标签生产与校验"""
print(f"[{datetime.now()}] 开始生产 {target_date} 的 RFM 标签…")

self.connect_db()

# 1. 执行 RFM SQL(由 Airflow 上游任务完成)
# 2. 数据量校验
today_cnt, deviation = self.check_data_volume(target_date)
print(f"今日标签用户数: {today_cnt}, 偏差率: {deviation:.2%}")

if deviation > 0.2: # 超过20%偏差
self.send_alert(
'数据量异常',
f'{target_date} 标签用户数偏差 {deviation:.1%},请检查上游数据!'
)
raise ValueError(f"数据量偏差过大: {deviation:.2%}")

# 3. 分布校验
distribution = self.check_distribution(target_date)
print("各分层用户分布:")
for tier, cnt, pct in distribution:
print(f" {tier}: {cnt} ({pct:.1f}%)")

# 4. 写入标签生产日志
log_sql = """
INSERT INTO rfm_tag_production_log
(dt, total_users, status, create_time)
VALUES (%s, %s, 'SUCCESS', NOW())
"""
with self.conn.cursor() as cursor:
cursor.execute(log_sql, (target_date, today_cnt))
self.conn.commit()

print(f"[{datetime.now()}] {target_date} RFM 标签生产完成!")
self.conn.close()

# ========== 使用示例 ==========
if __name__ == '__main__':
db_config = {
'host': 'your-mysql-host',
'port': 3306,
'user': 'data_platform',
'password': '***',
'database': 'user_profile'
}

pipeline = RFMTagPipeline(db_config)
# 通常由 Airflow 传入执行日期
target_date = '2026-07-22'
pipeline.run(target_date)

四、生产环境的稳定性保障

把 RFM 从一次性分析变成每日自动生产的标签系统,最大的挑战其实是稳定性。

第一个坑是数据延迟。RFM 依赖前 90 天的订单数据,但上游 ODS 表偶尔会延迟(比如大促期间写入量暴增导致延迟 2-3 小时)。我们的做法是设置一个"最晚等待时间":凌晨 3 点开始跑,如果上游数据还没到齐,就先用昨天产出的标签兜底。虽然时效性差了点,但至少不会让下游系统吃到空数据。

第二个坑是分位点抖动。RFM 的分位点每天重算,在大促、节假日等流量高峰,分位点会剧烈变化,导致大量用户的标签在一夜之间"降级"——这对运营策略影响很大。解决方案是引入滑动窗口平滑:用过去 7 天的分位点均值作为当天切割依据,大幅减少了标签抖动。

第三个坑是跨业务线口径不统一。比如电商业务算的是"下单金额",内容付费算的是"付费内容消费金额",会员系统算的又是"订阅续费金额"。如果各团队各算各的,数据口径就乱套了。我们的方案是在 DW 层统一"消费金额"的定义,各业务线在此基础上加定制字段,保证底层一致性。

五、总结

RFM 模型的核心思路几十年没变过,但把它做成一个稳定可靠的工程系统,需要处理的问题远比"用 Excel 算分档"复杂得多。从数据口径统一、到分位点平滑、到异常校验兜底,每一个细节都直接影响着下游业务能否正常运转。

自动化标签系统的价值在于"一致性"和"可追溯性"。任何人打开标签表,都能清楚地看到一个用户为什么被分到"重要挽回客户"——因为 R=5, F=2, M=5,有据可查。而不是"我觉得这个用户挺重要的"这种拍脑袋的判断。

如果你也在做用户分层相关的工作,欢迎评论区聊聊你的方案!

赞(0)
未经允许不得转载:171主机测评 » RFM 模型工程化落地:用户分层从 Excel 到自动化标签系统
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址