ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Angel 分布式 Swing 算法实战:基于用户行为计算 Item 相似度的推荐召回方案

Angel 分布式 Swing 算法实战:基于用户行为计算 Item 相似度的推荐召回方案 人工智能机器学习分布式训练图计算后端【免费下载链接】angelA Flexible and Powerful Parameter Server for large-scale machine learning项目地址https://gitcode.com/gh_mirrors/an/angel点击查看免费下载本文以 Angel 开源项目Spark on Angel中的 Swing 图算法为例系统讲解如何利用用户行为二部图计算 item 之间的相似度并输出itemId itemId score形式的相似度结果用于推荐召回。读者将掌握 Swing 的相似度公式与各参数alpha/beta/gamma/topFrom/topTo 等的真实含义、Spark on Angel 集群任务提交方式、内存资源估算方法以及常见超时问题的排查思路。文中涉及的公式、参数与执行流程均可在仓库源码中找到对应实现。1. 算法介绍为什么共同购买者越少Item 越相似Swing 是一种利用用户行为如购买、点击计算 item 与 item 之间相似度的算法通常用于推荐系统的召回阶段。其基本原理是如果两个 item 的购买用户集合共同关联的 user越少那么这些 item 之间的相似度越高。也就是说一对 item 被少数几个口味一致的忠实用户共同购买比被大量泛化用户共同购买更能说明它们之间的相似关系——后者往往来自用户的普适偏好区分度低。在二部图中记Ui为购买过 item_i 的用户集合Uj为购买过 item_j 的用户集合Iu为用户 u 购买的 item 集合Iv为用户 v 购买的 item 集合Swing 相似度计算公式如下公式包含两部分相似度累加项sim(i,j) Σ W / (alpha |Iu ∩ Iv|)对同时购买过 item_i 与 item_j 的用户对 (u, v) 求和。分母中的|Iu ∩ Iv|是用户 u 与用户 v 共同购买过的 item 数量公共购买越多单对用户的贡献越小权重项W (|Iu| beta)^gamma × (|Iv| beta)^gamma其中gamma通常取值范围为[-1, 0]英文文档记为[-1, 0)beta与gamma共同构成对两条购买记录长度的惩罚——用户购买列表越长其口味信号越弱贡献权重越低。上述公式在源码 SwingOperator.scala 中实现为def Score(Iu:Array[Long], Iv:Array[Long], alpha:Float, weight: Float 1.0f): Float { weight / (alpha ArrayUtils.intersectCount(Iu, Iv)) }而权重wj (|Iu| beta)^gamma × (|Iv| beta)^gamma由swing()方法中的math.pow(uItems.length beta, gamma)与math.pow(vItems.length beta, gamma)逐对计算见 SwingOperator.scala。可以看到alpha作为分母平滑项存在beta对购买列表长度做偏移gamma为负值时即对长购买记录降权。2. 输入输出与数据格式Swing 的输入是一张无权的 user-item 二部图输出是 item 之间的相似度结果两者均为 HDFS 上的文本文件参数含义说明input输入二部图HDFS 路径不带权每行表示一条边格式为userId itemId由sep指定的分隔符分开output输出HDFS 路径每行表示一对 item 及其相似度格式为itemId itemId scoresep分隔符输入中每条边的起始顶点、目标顶点之间的分隔符如tab、空格等示例输入10001 20001 10001 20002 10002 20001 10002 20003示例输出score为相似度itemI itemJ与itemJ itemI都会写出对称结果20001 20002 0.125 20002 20001 0.125在源码中输入由 SwingExample.scala 通过GraphIO.load(input, isWeighted false, srcIndex, dstIndex, sep)加载输出通过GraphIO.save(mapping, output)保存。srcIndex/dstIndex可指定边文件中源点user与终点item所在的列索引默认均为 0、1。3. 算法参数详解默认值、取值范围与调优含义3.1 核心算法参数参数含义默认值topFrom将 item 按出现次数即被多少用户购买由大到小排序后仅计算排名在[topFrom, topTo)之间的 item 与其他所有 item 的相似度适用于只计算尾部长尾商品之间的相似度0topTo参见topFrom解释构成排名区间上界0alpha公式中的平滑项对应源码中的 delta 参数0beta公式中对购买记录长度的偏移惩罚参数5gamma公式中对购买记录长度的幂次惩罚参数取值范围[-1, 0]-0.3partitionNum数据分区数Spark RDD 数据的分区数量1psPartitionNum参数服务器PS上模型的分区数量取spark.ps.instancesuseBalancePartition参数服务器对输入数据节点存储划分是否采用均衡分区。如果输入节点的索引分布不均匀建议选择 true—storageLevelRDD 存储级别可选DISK_ONLY/MEMORY_ONLY/MEMORY_AND_DISKMEMORY_ONLY参数默认值在仓库中有多处印证HasBeta.scala 中setDefault(beta, 5)HasGamma.scala 中setDefault(gamma, -0.3f)示例入口 SwingExample.scala 中alpha默认0、beta默认5、gamma默认-0.3、storageLevel默认MEMORY_ONLY。调优建议由公式与文档推断gamma越接近 -1对购买列表长的用户惩罚越重相似度越倾向于由短列表、高忠诚度用户贡献topFrom/topTo组合适合关注长尾 item 的相似度例如新上架商品召回当区间覆盖 item 数超过 1 亿时源码 Swing.scala 会打印提示信息storageLevel用于控制 user-item 邻接表 RDD 的持久化级别内存紧张时可切到DISK_ONLY换取稳定性。3.2 示例独有的进阶参数集群示例 SwingExample.scala 还解析了以下进阶参数文档未逐一列出但同样可通过命令行传入参数含义示例默认值batchSize向 PS 初始化邻接表时每个批次处理的节点数10000pullBatchSize计算阶段从 PS 批量拉取邻居的分批大小200superItemThreshold共同用户数超过该阈值的 item 对会进入超级 item 对单独批处理通道2000superItemPairBatch超级 item 对的批处理大小800srcIndex/dstIndex输入文件中源点、终点的列索引0 / 1cpDirSpark checkpoint 目录默认取GraphIO.defaultCheckpointDir其中superItemThreshold与superItemPairBatch对应 Swing.scala 中的超级 item 对处理逻辑当某 item 对的共同用户数超过阈值时直接按(item_i, item, -1f)标记再统一交给calcSuperItemPairs以批处理方式补算分数避免单个分区内出现超大规模集合运算。4. 资源估算与内存配置Swing 在 PS 上存储的是user-item 邻接表每个用户节点挂其购买过的 item 列表内存占用可以按边数估算。文档给出了明确的内存配比建议PS 侧ps.instance与ps.memory的乘积是 PS 总的配置内存。为保证 Angel 不因内存不足而挂掉需要配置约为 PS 上数据存储量2 倍左右的内存Spark Executor 侧num-executors与executor-memory的乘积是 executors 总的配置内存最好能存下2 倍的输入数据。如果内存紧张1 倍也可以接受但运行会相对慢一些。文档给出的估算示例100 亿条边集的输入大约有 160G 大小20G × 20即 20 个 executor、每个 10G 内存左右的配置是足够的在资源实在紧张的情况下尝试加大分区数目partitionNum来分摊压力。这一内存模型与源码的执行流程一致edges先被persist(StorageLevel.DISK_ONLY)落盘见 Swing.scala随后userItemNeighborTable按storageLevel持久化并 push 到 PS计算阶段每个 partition 再从 PS 按pullBatchSize批量拉取邻居——PS 上的邻接表大小直接决定 ps.memory 需求。5. 任务提交示例Spark on Angel5.1 集群模式提交以下为文档提供的完整 Yarn 集群提交示例参数通过key:value形式跟在 jar 之后传入inputhdfs://my-hdfs/data outputhdfs://my-hdfs/output source ./spark-on-angel-env.sh $SPARK_HOME/bin/spark-submit \ --master yarn-cluster \ --conf spark.ps.instances1 \ --conf spark.ps.cores1 \ --conf spark.ps.jars$SONA_ANGEL_JARS \ --conf spark.ps.memory10g \ --name swing angel \ --jars $SONA_SPARK_JARS \ --driver-memory 5g \ --num-executors 1 \ --executor-cores 4 \ --executor-memory 10g \ --class org.apache.spark.angel.examples.graph.SwingExample \ ../lib/spark-on-angel-examples-3.3.0.jar input:$input output:$output sep:tab storageLevel:MEMORY_ONLY useBalancePartition:true \ partitionNum:4 psPartitionNum:1其中spark.ps.*系列配置用于启动 Angel PS参数服务器spark.ps.instances1表示 1 个 PS 实例、spark.ps.cores1分配 1 个核、spark.ps.memory10g分配 10G 内存$SONA_ANGEL_JARS/$SONA_SPARK_JARS由spark-on-angel-env.sh环境脚本注入对应 Angel 与 Spark 的依赖 jar。当前仓库中与文档同名入口对应的示例源码位于 cluster/SwingExample.scala包名为com.tencent.angel.spark.examples.cluster.SwingExample任务提交的通用环境准备可参考 spark-on-angel/README.md 与 spark_on_angel_quick_start.md。5.2 本地模式运行仓库还提供了本地模式示例 local/SwingExample.scala通过mode参数默认yarn-cluster切换运行方式便于在小数据上快速验证# mode 可指定为 local 等本地 master --mode local --input 本地或HDFS输入 --output 输出路径 sep:tab storageLevel:MEMORY_ONLY partitionNum:4 psPartitionNum:16. 分布式执行流程从二部图到相似度结果结合 Swing.scala 的transform方法一次 Swing 任务在分布式环境下的完整流程如下读取边数据NeighborDataOps.loadEdges读取 user-item 边并落盘DISK_ONLY打印抽样结果可选筛选目标 item 区间当topTo topFrom时对 item 按出现次数降序排序并取排名[topFrom, topTo)的 item 集合见 Swing.scala构建并统计邻接表SwingOperator.userItem2NeighborTable将边按 key 分组为(节点, 邻居数组)其中邻居数组做distinct.sorted去重排序见 SwingOperator.scalastats方法统计出 min/max id、节点数、边数与最大/最小度数初始化 PS 模型PSContext.getOrCreate启动 PS创建SimpleNeighborTableModel按batchSize分批将 user-item 邻接表 push 到 PS 并checkpoint()分区并行计算相似度itemUserNeighborTable.mapPartitionsWithIndex在每个分区内按pullBatchSize从 PS 批量拉取用户邻居再调用swing()计算 item 对分数见 SwingOperator.scala超级 item 对补算共同用户数超过superItemThreshold的 item 对进入calcSuperItemPairs批处理通道聚合后与主结果union输出过滤掉标记为-1f的中间结果按(itemI, itemJ, score)三元组写出 DataFrame见 Swing.scala。值得注意的实现细节swing()对每对用户计算权重wj时使用了math.pow(uItems.length beta, gamma)Scala 的math.pow返回 Double并与交集计数结合得到wj / (alpha intersectCount(Iu, Iv))对称的 item 对(item, item_i)也会被写出保证输出结果对每对 item 双向可见便于直接用于召回索引。7. 常见问题排查文档记录了一个典型的高频故障任务运行约 10 分钟时挂掉。可能原因Angel 申请不到资源。由于该任务基于 Spark on Angel 开发实际涉及 Spark 和 Angel 两个系统向 Yarn 申请资源时是独立进行的——Spark 任务拉起之后由 Spark 向 Yarn 提交 Angel 的任务如果不能在给定时间内申请到资源就会报超时错误导致任务挂掉解决方案确认资源池有足够的资源添加 Spark confspark.hadoop.angel.am.appstate.timeout.msxxx调大超时时间默认值为600000也就是 10 分钟。此外结合第 4 节的资源估算若任务频繁 OOM 或执行缓慢应优先检查 PS 内存ps.instance × ps.memory是否达到邻接表存储量的 2 倍、executor 总内存是否能容纳 2 倍输入数据并在内存紧张时增大partitionNum分区数。小结Swing 算法以共同购买者越少越相似为核心思想通过 alpha/beta/gamma 三个参数实现对购买记录长度的平滑与惩罚特别适合长尾 item 的相似度召回。在 Angel 的 Spark on Angel 框架下Swing 将邻接表存放在 PS 上、以分区并行方式批量拉取计算兼顾了超大二部图的可扩展性。结合本文给出的参数默认值、内存估算公式与提交脚本读者可以直接在 Yarn 集群上跑通input - output的完整链路。更多细节可进一步阅读仓库中的 swing_en.md 英文文档以及 Swing.scala 源码实现。赞分享人工智能机器学习分布式训练图计算后端【免费下载链接】angelA Flexible and Powerful Parameter Server for large-scale machine learning项目地址https://gitcode.com/gh_mirrors/an/angel点击查看免费下载相关推荐基于 Angel 图计算框架的 Swing 推荐召回算法实现与实战指南基于 Angel 图计算框架的 Swing 推荐召回算法实现与实战指南 导读 本文讲解 Angel 项目中基于 Spark On Angel 实现的 Swing人工智能机器学习分布式训练图计算后端Angel 图算法实战基于 Spark On Angel 的 BruteForce 暴力最近邻TopK 相似度搜索详解Angel 图算法实战基于 Spark On Angel 的 BruteForce 暴力最近邻TopK 相似度搜索详解 导读 BruteForce 是 A人工智能机器学习分布式训练图计算后端Angel 图计算系列基于 Spark On Angel 的分布式 Closeness 接近中心性算法实践Angel 图计算系列基于 Spark On Angel 的分布式 Closeness 接近中心性算法实践 本文聚焦 Angel 开源仓库中 docs/alg人工智能机器学习分布式训练图计算后端上一篇3分钟掌握百度文库文档免费保存技巧一键获取干净PDF下一篇OBS背景移除插件从AI算法到专业直播的完整技术解构创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表