拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于Spark的亿级用户聚类分析实战:K-Means客户细分全流程

很多做数据分析和用户增长的朋友一聊到客户细分第一反应就是用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, onuser_id, howouter)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(inputColsfeature_columns, outputColfeatures_raw) df_vector assembler.transform(df_user) # 标准化 scaler StandardScaler(inputColfeatures_raw, outputColfeatures, withStdTrue, withMeanTrue) scaler_model scaler.fit(df_vector) df_scaled scaler_model.transform(df_vector) # 选择K值跑了K3到K10记录轮廓系数 evaluator ClusteringEvaluator(featuresColfeatures, metricNamesilhouette) for k in range(3, 11): kmeans KMeans(featuresColfeatures, kk, seed42, maxIter30) model kmeans.fit(df_scaled) predictions model.transform(df_scaled) score evaluator.evaluate(predictions) print(fK{k}, silhouette_score{score:.4f})从运行日志来看K5的时候轮廓系数是0.412K6是0.398K4是0.376。轮廓系数整体不算很高但要注意这是在用户行为数据上很正常的现象用户群体之间的边界本来就模糊不像图像数据那样天然分得开。最终综合考虑业务解释性我选了K5把用户分成了五类每一类的画像特征都很清晰。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值环节参与了评审结合他们对用户的实际感知来判断聚类结果是否符合常识。另一个建议是如果条件允许可以把聚类的结果做成一个周期性的离线任务每周或者每月跑一次然后把用户的簇标签变化情况记录下来。有些用户这个月在“沉睡高价值”群体下个月突然变成了“高活跃高价值”这类变化本身就是很有价值的市场信号可以用来做触达策略。如果后续数据量继续增大还可以引入流式计算实时计算用户特征并更新聚类标签不过那就是另一个更复杂的项目了。最后如果对聚类特征工程这块想深入了解可以去看看网上那些电商用户消费行为聚类分析的案例理解了别人是怎么设计特征的自己动手做的时候会少走很多弯路。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门