很多做数据分析和用户增长的朋友,一聊到客户细分,第一反应就是用SQL跑几个RFM指标,然后手动分一下层。这种做法在数据量小、维度少的时候还行,可一旦用户量到了千万级,特征维度扩展到十几个的时候,传统方式基本就跑不动了。我之前在一家电商平台做过一次基于Spark的聚类分析项目,目的就是解决亿级用户的精细化分层问题,用到的核心算法是K-Means,同时也对比了Bisecting K-Means和高斯混合模型。这篇博文就是那次项目的完整复盘,从方案设计到代码实现再到调参经验,一次讲清楚,适合正在做用户画像、增长策略或者刚接触Spark MLlib的工程师参考。
1. 业务场景与方案设计
1.1 为什么客户细分要用聚类分析
客户细分这件事,本质上是把一群行为特征相似的用户归到同一个组里,然后针对不同组做差异化的运营策略。过去常见的做法是业务方凭经验定规则,比如“30天内下单超过5次的是高价值用户”“近7天未访问的是流失风险用户”之类的。这种规则虽然解释性强,但有一个很大的问题,就是你得先知道用户有哪些类型,才能定义出这些规则。当用户行为模式变得复杂,比如有些用户购买频次低但客单价极高,有些用户频繁浏览但从不下单,这些规则就很难覆盖全。
聚类分析解决的是"不知道有什么类型"的问题。它通过计算用户在各个特征维度上的距离,自动把相似的用户聚在一起,不需要提前标注标签,属于无监督学习。你只需要把特征准备好,剩下的群组划分交给算法完成。Spark MLlib里的聚类算法实现,刚好能支撑海量用户数据的计算需求。
1.2 为什么选Spark而不是单机Python
在项目选型时,我其实纠结过是用单机Python的scikit-learn还是Spark。后来测了一下数据量,用户特征表大概有1.2亿行,特征维度有26个,如果把这堆数据拉到一台机器上跑K-Means,内存就直接爆掉了,更别提还要做特征工程和多次迭代调参。Spark的核心优势在于它是分布式内存计算框架,能把数据分片到多台机器上并行计算,K-Means这种迭代型算法在Spark上天然合适,因为每轮迭代只需要在Driver端汇总一下聚类中心,再广播回各个Executor做下一轮。
另一个需要考虑的因素是项目里还有其他ETL任务,比如用户行为日志的清洗、订单数据的汇总加工,这些本来就在Spark集群上跑。把聚类分析放在同一个Spark平台上,能省掉大量数据搬运的时间。数据从Hive或者Parquet文件里读出来,算完直接写回表,整个链路是通的,不需要把数据导出来再导进去。
2. 数据集与特征工程
2.1 特征设计的核心思路:RFM模型的扩展
客户细分最经典的底层模型就是RFM,即最近一次消费时间(Recency)、消费频率(Frequency)和消费金额(Monetary)。这个模型的逻辑很简单:最近买过、经常买、买得多的人,价值自然更高。但实际操作中,只有这三个维度远远不够,我在这基础上扩展出了几类特征。
一是行为活跃度特征,比如近30天登录次数、近30天浏览商品数、平均每次会话时长。二是品类偏好特征,比如用户在美妆类目的购买占比、在数码类目的浏览占比。三是消费稳定性特征,比如消费间隔的变异系数、月度消费金额的波动情况。这些特征加在一起,能从多个角度刻画用户的真实状态。比如有两个用户RFM三个值完全一样,但一个只买打折商品,一个原价购买为主,两者价值其实是不同的,只有把价格敏感度这类特征加进去才能区分开。
2.2 数据清洗与特征变换的实际操作
特征工程这一步最耗时间,也最影响最终效果。原始数据来自好几张表,包括订单表、访问日志表、用户注册表,所以第一步是先把这些表关联起来。订单表按用户维度做聚合,计算出总消费金额、总订单数、最近下单时间等;访问日志表按用户聚合出浏览行为指标;最后再把这些结果按用户ID关联成一张宽表。
这一步有几个坑必须提前处理。第一是空值问题,比如一个用户注册了但从未下过单,那么订单金额字段就是空的,需要填充为0,但如果所有空值都用0填充,又会让“未下单”和“下单后退款”的用户混淆,所以最好单独加一个标志字段。第二是极值问题,消费金额的分布极度右偏,少数头部用户可能占了大部分销售额,如果不做处理,聚类结果会被这几个头部用户带偏。我用的处理方式是把金额取对数,也就是log1p变换,拉近数值之间的距离。
数据标准化也是必须做的一步,而且是很多新手容易忽略的。K-Means是基于欧氏距离计算的算法,如果某个特征的数值范围特别大,比如消费金额从0到几万,而登录次数只有0到几十,那么距离计算会被消费金额完全主导,其他特征就失去了意义。Spark MLlib里的StandardScaler可以把每个特征变换成均值为0、方差为1的标准分布,这是聚类前必须做的一个步骤。
3. 聚类算法选型与核心参数
3.1 K-Means算法原理与Spark实现特点
K-Means是聚类算法里应用最广泛的,思想也最直观。算法先随机选K个点作为初始聚类中心,然后把每个样本点分配到距离它最近的中心所在的簇,接着重新计算每个簇的中心点,再分配、再更新,迭代直到中心点不再变化或者达到预设的迭代上限。整个过程在数学上是在最小化每个样本点与所属簇中心的欧氏距离平方和。
Spark的MLlib里对K-Means做了分布式优化。在分布式环境下,每次迭代时每个Executor节点上会计算自己分到的数据与当前聚类中心的距离,并汇总局部统计量,然后在Driver端更新全局聚类中心。Spark的实现还加了初始化优化机制,用的是K-Means||算法,这个算法不是单纯随机选初始点,而是通过多次采样来选取分布更均匀的初始中心,能明显减少后续迭代次数,同时降低落到局部最优解的概率。
我在用的时候,K-Means的K值并不是直接拍脑袋定的,而是结合了业务可解释性和算法评估指标来选的。算法层面上用到了轮廓系数,它衡量的是样本点与自身簇内样本的紧密程度,以及与相邻簇样本的分离程度,取值范围在-1到1之间,越大说明聚类效果越好。我分别试了K值从3到10,把每个K值的轮廓系数打出来对比,同时把每个类别的样本量和业务特征打印出来给运营同学看,两边都满意才定下来。
3.2 不同聚类算法的对比与适用场景
除了K-Means,项目里我还试过Bisecting K-Means和高斯混合模型(GMM),这里也顺手做个对比。K-Means假设每个簇是凸形的,就是圆球形分布,而且强制每个样本点只能属于一个簇,这可能跟现实不完全匹配。Bisecting K-Means是K-Means的层级版本,先所有数据当成一个簇,然后逐步二分,每次选当前最大的簇继续分裂,算法效率更高,而且在某些非球形分布的数据上表现更好。
高斯混合模型则更灵活一些,它允许一个样本点以概率的形式属于多个簇,能更好地处理簇与簇之间边界模糊的情况。但GMM的计算复杂度比K-Means高不少,而且在数据量特别大的时候迭代速度更慢。我当时跑了一版GMM在同样数据上,时间大概是K-Means的三倍,效果提升并不明显,所以最终上线还是用了K-Means。
如果数据形状比较不规则,还可以考虑DBSCAN这种基于密度的算法。网上也有一些"基于K-Means与DBSCAN的电商用户消费行为聚类分析"的案例研究,思路很好,但DBSCAN在Spark上的实现不太完善,调参难度大,要知道近邻半径和最少样本数这两个参数在分布式环境下的调试成本很高,我这次项目就没有选择它。
4. 基于Spark的完整实操流程
4.1 Spark环境与集群配置要点
在动手写代码之前,集群环境要先准备好。如果公司内部已经有大数据的Hadoop集群,那Spark直接部署在YARN上就行,资源可以动态分配。我这次用的是独立的Spark集群,三台worker节点,每台配了64GB内存、16核CPU。跑亿级数据的K-Means,资源规划上要注意给Driver端留够内存,因为K-Means的聚类中心虽然不大,但广播变量和累加器的数据还是会占一些内存,而且Driver还要处理任务调度带来的开销。
Spark的内存配置有几个关键参数要注意。spark.executor.memory决定了每个Executor能用的堆内存大小,spark.executor.cores决定Executor的并行度。我建议把Executor内存分配给存储和执行两部分的比例纳入考虑,存储部分用来缓存DataFrame和RDD,执行部分用来跑Shuffle和聚合计算。如果比例设置不当,容易出现内存溢出或者频繁的GC,尤其在做特征工程阶段有大量宽表关联操作,这一步内存压力还是很大的。
4.2 特征表构建与DataFrame操作
数据源从Hive里读出来之后,先用Spark SQL做预处理。下面是我实际跑过的代码片段,为了方便展示,做了脱敏和简化处理。
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.clustering import KMeans from pyspark.ml.evaluation import ClusteringEvaluator spark = SparkSession.builder .appName("customer_segmentation") .enableHiveSupport() .getOrCreate() # 读取订单表和用户行为表,按用户维度聚合特征 df_orders = spark.sql(""" SELECT user_id, COUNT(DISTINCT order_id) AS order_cnt, SUM(order_amount) AS total_amount, DATEDIFF(CURRENT_DATE, MAX(order_time)) AS recency_days, AVG(order_amount) AS avg_order_amount FROM dwd_order_info WHERE dt = '2024-11-30' GROUP BY user_id """) df_behavior = spark.sql(""" SELECT user_id, COUNT(*) AS visit_cnt, COUNT(DISTINCT product_id) AS browse_product_cnt, SUM(session_duration) AS total_duration FROM dwd_user_behavior_log WHERE dt = '2024-11-30' GROUP BY user_id """) # 关联成宽表 df_user = df_orders.join(df_behavior, on="user_id", how="outer")4.3 K-Means聚类训练的完整过程与调参记录
宽表构建完成后,是特征选择与向量装配环节。这一步需要用VectorAssembler把多个特征列合成一个向量列,然后做标准化,最后灌进K-Means模型里训练。
feature_columns = [ "recency_days", "order_cnt", "total_amount", "avg_order_amount", "visit_cnt", "browse_product_cnt", "total_duration" ] # log变换处理金额类和时长类的长尾分布 for col in ["total_amount", "avg_order_amount", "total_duration"]: df_user = df_user.withColumn(col, F.log1p(F.col(col))) # 缺失值填充 df_user = df_user.fillna(0) # 组装特征向量 assembler = VectorAssembler(inputCols=feature_columns, outputCol="features_raw") df_vector = assembler.transform(df_user) # 标准化 scaler = StandardScaler(inputCol="features_raw", outputCol="features", withStd=True, withMean=True) scaler_model = scaler.fit(df_vector) df_scaled = scaler_model.transform(df_vector) # 选择K值:跑了K=3到K=10,记录轮廓系数 evaluator = ClusteringEvaluator(featuresCol="features", metricName="silhouette") for k in range(3, 11): kmeans = KMeans(featuresCol="features", k=k, seed=42, maxIter=30) model = kmeans.fit(df_scaled) predictions = model.transform(df_scaled) score = evaluator.evaluate(predictions) print(f"K={k}, silhouette_score={score:.4f}")从运行日志来看,K=5的时候轮廓系数是0.412,K=6是0.398,K=4是0.376。轮廓系数整体不算很高,但要注意这是在用户行为数据上很正常的现象,用户群体之间的边界本来就模糊,不像图像数据那样天然分得开。最终综合考虑业务解释性,我选了K=5,把用户分成了五类,每一类的画像特征都很清晰。
4.4 聚类结果解读与业务落地的映射方式
模型跑完之后,最关键的一步是把聚类标签对应回原始用户,并做群体画像分析。这一步是在聚类模型输出的预测结果上,重新按簇分组计算各维度的均值和中位数,整理成一张用户分群特征表。
我当时得到五个群体的轮廓大概是这样的:
- 群体0:高活跃高价值用户,占比约8%,贡献了超过45%的GMV,平均客单价高,复购频次高,是核心种子用户
- 群体1:中活跃中价值用户,占比约22%,消费频次稳定,处于成长阶段,适合做品类拓展和交叉销售
- 群体2:高活跃低价值用户,占比约15%,访问频繁但消费很少,有转化潜力,适合发优惠券激活
- 群体3:低活跃中高价值用户,占比约18%,消费金额不低但已经不常来了,属于沉睡高价值用户,需要做召回
- 群体4:低活跃低价值用户,占比约37%,基本处于流失状态,维护成本高,不适合投入大量运营资源
得到了这些分群结果,我还同时生成了每个用户群的用户ID清单,写入到一个Hive表里,下游的运营系统在推送消息、设计活动的时候,直接按标签筛选人群就可以了,这样整个链路就跑通了。
5. 常见问题与排查技实录
5.1 聚类结果为空簇或者簇大小失衡
我在前期调试的时候就遇到过跑完K-Means后,某个K值下有一个簇只有几百个用户,而其他簇有上千万用户的情况。这个问题的根源多半出在初始聚类中心的选择上,如果初始中心落在了离群点附近,某个簇在迭代过程中就可能被慢慢掏空。Spark的K-Means实现了K-Means||初始化,已经比朴素随机初始化稳定很多,但极端情况下还是可能碰到局部最优解。
解决思路是换随机种子多跑几次,然后把每次运行结果的簇大小打出来对比,选择一个稳定的结果。代码里我给KMeans设置了多个不同的seed值,比如42、2023、7,对比三组结果,选簇大小分布更均匀且轮廓系数更高的那一个作为最终模型。
5.2 特征量纲不统一导致聚类结果被单一特征主导
这是一个特别容易踩的坑,我也没少往里掉。最开始我偷懒,做完log变换以后没有做标准化,直接把向量喂给了K-Means,结果聚类出来的群组基本就是按消费金额一个维度分的,其他特征在里面毫无存在感。后来加上StandardScaler之后,各特征的贡献度才均衡起来。
判断是否存在这个问题的办法很简单,训练完之后把每个簇的中心点向量打印出来,看一下各维度的数值分布。如果某个维度的数值跨度远大于其他维度,那多半就是标准化出问题了。另外要说的一点是StandardScaler要在大样本上fit,样本量小的时候均值和标准差的估计不稳定,也会影响效果。
5.3 迭代不收敛或者训练时间过长
在数据量大、特征维度高的情况下,K-Means的迭代收敛速度会明显下降。解决办法有几个维度。一是调大maxIter的同时设置tol参数,就是聚类中心变化的容忍度,中心变化小于这个值时提前终止迭代。二是检查数据分区数,分区太少容易导致并行度不够,分区太多又会产生大量调度开销。我当时把数据重分区到每个分区大约500MB到1GB大小,训练速度有明显提升。三是检查是否存在严重的数据倾斜,某些热点用户相关的数据量特别大,会导致个别任务耗时过长,这时候需要按用户ID加盐重新分区,把热点数据打散。
6. 项目心得体会与扩展建议
这个项目从开发到上线,前后大概用了两周时间。如果说有什么经验是我想特别强调的,那就是在开始写代码之前,一定要先把业务问题想清楚。聚类分析是一个无监督的过程,模型本身并不知道用户价值高还是低,它只能帮你把行为模式相近的人放到一起。最终的判断和解释,还是要靠业务方一起参与。我在项目过程中就让运营同事在确定K值环节参与了评审,结合他们对用户的实际感知来判断聚类结果是否符合常识。
另一个建议是,如果条件允许,可以把聚类的结果做成一个周期性的离线任务,每周或者每月跑一次,然后把用户的簇标签变化情况记录下来。有些用户这个月在“沉睡高价值”群体,下个月突然变成了“高活跃高价值”,这类变化本身就是很有价值的市场信号,可以用来做触达策略。如果后续数据量继续增大,还可以引入流式计算,实时计算用户特征并更新聚类标签,不过那就是另一个更复杂的项目了。
最后,如果对聚类特征工程这块想深入了解,可以去看看网上那些"电商用户消费行为聚类分析"的案例,理解了别人是怎么设计特征的,自己动手做的时候会少走很多弯路。