
推荐系统中的特征工程流水线从离线计算到在线服务的架构设计一、推荐系统特征工程的架构分层推荐系统的特征工程与传统的机器学习特征工程有本质差异它不仅需要处理大规模、多源、异构的数据还必须在离线训练和在线推理两个环境中保持特征计算逻辑的一致性。训练时在Spark上用Python计算的特征推理时需要在低延迟的Go/Java服务中复现——任何微小的实现差异都可能导致训练-推理偏差Training-Serving Skew直接损害模型的线上效果。从架构视角看推荐系统的特征工程可以划分为四个层次特征定义层特征的语义描述和元数据管理、离线计算层大规模批处理特征生成、在线计算层低延迟特征实时生成、特征存储层特征的持久化与低延迟读取。各层之间通过统一的特征注册中心来保证语义一致性。二、离线特征计算的工程范式离线特征计算负责生成推荐系统中体量最大的特征类别——聚合统计类特征用户过去N天的点击率、商品的7日曝光转化率等和序列特征用户最近K次交互的商品ID序列。这些特征的计算通常依赖数天甚至数月的行为日志数据量级在TB-PB级别必须依赖分布式计算框架。Spark SQL是实现离线特征计算的主流工具。以下模式代表了典型的用户行为聚合特征的计算范式 使用PySpark计算推荐系统常用的用户行为聚合特征 from pyspark.sql import SparkSession, Window from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, LongType, FloatType spark SparkSession.builder \ .appName(FeatureEngineering) \ .config(spark.sql.adaptive.enabled, true) \ .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ .getOrCreate() # ---- 假设的行为日志表结构 ---- # user_id: 用户ID, item_id: 商品ID, action: 行为类型(click/like/buy) # event_time: 事件时间戳, price: 商品价格 def compute_user_aggregation_features( behavior_table: str, window_days: list[int] [1, 7, 30] ) - DataFrame: 计算用户聚合特征多时间窗口的行为统计。 为每个时间窗口生成一组特征 - action_count: 行为总数 - action_type_distribution: 各行为类型的比例 - distinct_items: 去重商品数 - avg_price: 平均浏览价格 Args: behavior_table: Hive行为日志表名 window_days: 滑动窗口的尺寸列表天 Returns: DataFrame: 包含所有窗口特征的宽表每行对应一个user_id # 读取指定时间范围的行为数据 max_window max(window_days) current_ts F.current_timestamp() df spark.table(behavior_table).filter( F.col(event_time) F.date_sub(current_ts, max_window) ) # 为每个窗口分别计算特征然后join result_df None for days in window_days: window_df df.filter( F.col(event_time) F.date_sub(current_ts, days) ) # 用户级聚合 agg_df window_df.groupBy(user_id).agg( F.count(action).alias(faction_count_{days}d), # 各行为类型的计数 F.sum(F.when(F.col(action) click, 1).otherwise(0)).alias(fclick_count_{days}d), F.sum(F.when(F.col(action) like, 1).otherwise(0)).alias(flike_count_{days}d), F.sum(F.when(F.col(action) buy, 1).otherwise(0)).alias(fbuy_count_{days}d), # 去重商品数 F.countDistinct(item_id).alias(fdistinct_items_{days}d), # 平均浏览价格 F.avg(price).alias(favg_price_{days}d), ) # 计算行为比例特征 agg_df agg_df.withColumn( fclick_ratio_{days}d, F.col(fclick_count_{days}d) / F.col(faction_count_{days}d) ).withColumn( fbuy_conversion_{days}d, F.col(fbuy_count_{days}d) / F.col(fclick_count_{days}d) ) if result_df is None: result_df agg_df else: result_df result_df.join(agg_df, user_id, outer) return result_df def compute_item_sequence_features( behavior_table: str, sequence_length: int 20 ) - DataFrame: 计算用户最近交互的商品序列特征。 按时间排序取最近N次交互的商品ID列表 可用于序列推荐模型如SASRec、BST。 Args: behavior_table: 行为日志表 sequence_length: 序列长度 Returns: DataFrame: user_id item_sequence数组列 df spark.table(behavior_table) # 按用户分区、按时间排序取最近N条 window_spec Window.partitionBy(user_id).orderBy(F.col(event_time).desc()) sequence_df df.withColumn(rank, F.row_number().over(window_spec)) \ .filter(F.col(rank) sequence_length) \ .groupBy(user_id) \ .agg( # collect_list按rank排序收集item_id F.collect_list(F.struct(rank, item_id)).alias(item_seq_struct) ) \ .withColumn( item_sequence, F.col(item_seq_struct.item_id) ) \ .drop(item_seq_struct) return sequence_df三、在线特征计算的延迟约束与缓存策略在线推理环节推荐系统需要在100ms内完成特征拉取、模型计算和排序。在这个延迟预算中特征获取通常占据40-60ms——是从离线批处理到在线服务的最大瓶颈。特征在线计算面临的核心问题是哪些特征应该离线预计算存入KV存储哪些特征必须在请求到达时实时计算决策矩阵如下用户长期统计特征30天点击率、历史购买均价离线预计算存入Redis请求时O(1)读取。更新频率为天级别T1。用户短期行为特征最近5次点击、当前会话内的浏览序列实时计算。这类特征时效性敏感T1更新会导致推荐滞后于用户当前兴趣。通常在API服务的内存中维护用户最近的会话状态。上下文特征当前时间、设备类型、网络环境请求携带无需存储。四、训练-推理一致性保障机制训练-推理偏差是特征工程中最隐蔽的质量风险。它发生在训练时的特征计算逻辑与推理时不一致的情况下——典型的来源包括离线特征使用窗口结束时间作为参考点在线推理使用请求到达时间离线计算使用全量数据聚合在线查询可能因Redis分片导致部分数据缺失浮点数精度差异Spark的double vs Go的float64。保障一致性的工程实践包括第一特征计算逻辑代码化将特征的计算公式以配置文件形式管理训练和推理共用同一份特征配置由不同的运行时引擎Spark/Go各自解析。避免在两个系统中分别手写计算逻辑。第二离线在线一致性监控定期每小时抽取少量在线请求将其特征值与同一时刻的离线计算结果进行diff对比差异超过阈值如5%时触发告警。第三特征版本化管理每次修改特征计算逻辑时生成新的特征版本号。训练数据和在线特征库使用版本号关联确保模型训练所用的特征定义与在线服务完全一致。不支持特征版本的回溯修改。五、总结推荐系统的特征工程流水线是一个横跨离线批处理和在线实时服务的分布式系统工程。四个架构层次——特征定义、离线计算、在线计算、特征存储——需要通过统一的元数据管理和版本控制来维持一致性。核心的工程取舍发生在特征的新鲜度与计算延迟之间用户长期统计特征通过日级离线预计算Redis缓存实现毫秒级读取用户短期行为特征通过请求时实时聚合来捕获即时的兴趣变化。训练-推理偏差是最隐蔽但影响最大的质量问题应通过特征计算逻辑的统一配置化、定期的离在线数据diff监控、以及特征版本化追溯来系统性防控。