
1. 从零开始理解Spark推荐系统第一次接触推荐系统时我被各种算法名词绕得头晕。直到用Spark真正实现了一个电商推荐模块才发现原来核心思想这么简单。想象你是个书店老板要给学生推荐书籍通常有两种做法要么看看这个学生喜欢什么书找类似的书推荐基于物品要么看看和他相似的学生都喜欢什么书基于用户。这就是推荐系统最基础的两种思路。Spark的MLlib库把这些算法都封装好了我们只需要理解怎么用。先来看基于物品的协同过滤我用一个真实案例说明。假设我们有个用户评分数据文件ratingdata.txt格式是用户ID,物品ID,评分。用Spark处理这种数据特别方便from pyspark.mllib.recommendation import ALS, MatrixFactorizationModel, Rating # 加载数据 data sc.textFile(ratingdata.txt) ratings data.map(lambda l: l.split(,)).map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))实际项目中会遇到数据稀疏问题。比如有10万用户和1万商品评分矩阵中99%都是空缺。Spark的分布式矩阵计算能高效处理这种稀疏场景。我做过一个实验在单机Python上跑协同过滤100MB数据就内存溢出同样的算法用Spark处理10GB数据集群模式下只要几分钟。2. 进阶实战ALS算法与推荐优化交替最小二乘(ALS)算法是Spark推荐系统的王牌。它通过矩阵分解把用户和物品映射到潜在特征空间比如把书籍分解为科幻成分、文学价值等维度。我在电商项目中发现合理设置这三个参数最关键rank特征数通常10-200之间需要交叉验证iterations迭代次数5-20次足够收敛lambda正则化系数防止过拟合0.01是常用起点# ALS模型训练最佳实践 model ALS.train(ratings, rank50, iterations10, lambda_0.01)评估推荐质量时我习惯用RMSE和业务指标结合。曾遇到RMSE很好但实际推荐效果差的情况后来发现是测试集没有做时间划分——用未来数据预测过去当然不准。正确做法是按时间划分训练测试集或者用留一验证。线上部署要注意冷启动问题。我们的解决方案是对新用户先用热门商品推荐收集足够行为数据后再切到个性化推荐。Spark ML的Pipeline功能可以很方便地实现这种多阶段策略。3. 金融风控中的Spark机器学习金融场景对实时性要求更高。一次线上支付要在100ms内完成风控判断这对Spark流处理是挑战。我们的架构是用Spark Streaming预处理特征加载预训练好的随机森林模型进行实时预测。随机森林的特征重要性分析能帮我们理解风险因素。在贷款审批项目中发现借款期限和历史信用是最强预测因子。特征处理时要注意类别变量需要StringIndexer转换数值变量最好做分桶处理缺失值用中位数填充比直接删除更优from pyspark.ml.feature import VectorAssembler from pyspark.ml.classification import RandomForestClassifier assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) rf RandomForestClassifier(labelCollabel, featuresColfeatures, numTrees20)模型监控同样重要。我们设置了三道防线实时统计预测分布、定期回溯测试、人工抽样审核。当预测通过率异常波动时能立即报警。4. 深度学习在Spark中的实践文本分类是典型深度学习场景。处理垃圾短信识别时传统方法效果停滞在90%准确率改用多层感知器(MLP)后提升到97%。Spark的MLP实现虽然不如TensorFlow灵活但对结构化特征足够用。词向量处理是关键一步。我们发现向量维度100-300足够minCount设为5能过滤噪声词窗口大小5-10适合句子语境from pyspark.ml.classification import MultilayerPerceptronClassifier from pyspark.ml.feature import Word2Vec word2Vec Word2Vec(vectorSize100, minCount5, inputColmessage, outputColfeatures) layers [100, 50, 2] # 输入层、隐藏层、输出层 mlp MultilayerPerceptronClassifier(layerslayers, blockSize128, seed1234)部署时遇到内存问题——加载大模型导致Executor崩溃。解决方案是调整并行度和广播变量把模型拆分成多个小部分加载。现在我们的反欺诈系统每天处理千万级消息延迟控制在200ms内。