ARTICLE DETAIL

资讯详情

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

Feast Tiling 与中间表示(IR):流式时间窗口聚合的优化原理与实战配置

Feast Tiling 与中间表示(IR):流式时间窗口聚合的优化原理与实战配置 Feast Tiling 与中间表示IR流式时间窗口聚合的优化原理与实战配置【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast导读本文以 Feast 官方概念文档 docs/getting-started/concepts/tiling.md 为核心骨架深入讲解Tiling分片技术——一种针对流式时间窗口聚合的优化手段通过把数据预聚合进更小的时间片tile并保存中间表示Intermediate RepresentationsIR来保证合并结果的数学正确性。读完本文你将掌握为什么存平均值再合并是错的、avg/std/var需要保存哪些 IR、Feast 的 Sawtooth锯齿窗分片算法如何工作以及如何在StreamFeatureView中正确配置enable_tiling与tiling_hop_size并了解 Spark/Ray 计算引擎背后的源码实现。一、背景时间窗口聚合的两难困境在实时特征平台中最常见的需求之一是基于滚动时间窗计算特征例如过去 1 小时某用户的交易总额、平均值、标准差。当数据以每分钟甚至每秒的速率从 Kafka、Kinesis 或 PushSource 流入时如何高效、正确地维护这些滑动窗口聚合是所有流式计算面临的共同难题。传统方案通常只有两条路各有利弊每次从原始数据全量重算——每过一个窗口周期就把整个窗口内的数据重新扫描一遍。正确性没问题但代价极高窗口越长、事件越密集开销越大按片存储最终聚合值再合并——先把窗口切成若干小片每片算出最终值比如平均值然后尝试把各片的最终值合并成整个窗口的最终值。速度很快但数学上常常是错的。Tiling 正是为打破这一两难而生的第三方案按固定步长hop切出小时间片在片内计算并保存中间表示IR用 IR 而不是最终值参与合并从而同时获得增量计算的速度与数学合并的正确性。合并问题为什么平均值合并是错的下面的例子直观展示了直接合并最终值的错误WRONG: avg(tile1, tile2) ≠ (avg_tile1 avg_tile2) / 2 Example: tile1: [10, 20, 30] → avg 20 tile2: [100] → avg 100 正确合并的平均值: (102030100) / 4 40 错误合并的平均值: (20 100) / 2 60同样的合并陷阱也存在于标准差std/stddev方差var/variance中位数与百分位数任何需要知道全部取值的holistic整体性聚合这类聚合函数不满足可加性associativity——把分片的最终结果再求一次平均/再求一次标准差并不会还原全局结果。二、解决方案中间表示IR与聚合分类Feast 的答案是不存最终聚合值而是保存保留数学合并性质所需的中间数据。仍以平均值为例传统错误方式Tile 1: avg 20 Tile 2: avg 20 Merged avg (20 20) / 2 20 —— 错误使用 IR正确方式Tile 1: sum 60, count 3 Tile 2: sum 100, count 1 Merged: sum 160, count 4 Merged avg 160 / 4 40 —— 正确Feast 把所有聚合函数划分为两类对应完全不同的存储与合并策略。代数聚合Algebraic Aggregations这类聚合本身就是可合并的——直接把各片的同类型聚合值再套一次相同函数即可因此最终值就是 IR无需额外保存中间表示聚合函数存储值合并策略存储开销sumsumsum(tile_sums)1 列countcountsum(tile_counts)1 列maxmaxmax(tile_maxes)1 列minminmin(tile_mins)1 列源码中这一定义体现在 sdk/python/feast/aggregation/tiling/base.py 的get_ir_metadata_for_aggregation()当agg.function属于[sum, count, max, min]时直接返回IRMetadata(typealgebraic)不生成任何 IR 列。整体聚合Holistic Aggregations需要同时保存多个中间值合并时对 IR 做分量相减/相加再套用最终公式。平均值avg,mean存储 IRsum、count最终计算avg sum / count合并策略先对 sum 求和、对 count 求和再相除存储开销3 列最终值 2 个 IR对应源码中avg/mean分支生成的 IR 列名形如_tail_{feature_name}_sum与_tail_{feature_name}_countcomputation记为sum / count。标准差std,stddev存储 IRcount、sum、sum_of_squares最终计算variance (sum_sq - sum**2 / count) / (count - δ) std sqrt(variance) # δ 1 表示样本标准差sampleδ 0 表示总体标准差population合并策略对三个 IR 分别求和再套公式存储开销4 列最终值 3 个 IR方差var,variance存储 IRcount、sum、sum_of_squares最终计算与 std 相同但去掉sqrt()存储开销4 列最终值 3 个 IR一个重要的边界情况在源码中被显式拒绝count_distinct不支持 Tiling。get_ir_metadata_for_aggregation()会直接抛出ValueError提示改用enable_tilingFalse或选择代数聚合sum/count/min/max。三、Tiling 工作流程连续更新与增量计算Tiling 专为**频繁更新如每隔几分钟**的流式场景设计其核心思路是一次只算一个新增片其余片直接复用。1. 连续分片更新流程Stream Events → Partition by Hop Intervals → Compute IRs → Store Windowed Aggregations | | | | | | | └─ Online Store (Redis 等) | | └─ avg_sum, avg_count, std_sum_sq 等 | └─ 5 分钟 hop: [00:00-00:05], [00:05-00:10], ... └─ customer_id1: [txn1, txn2, txn3, ...]以 1 小时窗口、5 分钟 hop 为例整个窗口被切成 12 个片。每 5 分钟新事件到达只计算1 个新片5 分钟的数据复用内存中的11 个旧片最终聚合 合并这 12 个片1 新 11 复用。为什么快不用 Tiling每 5 分钟就要扫描整个 1 小时窗口1000 条事件使用 Tiling只处理 5 分钟的新事件旧片直接复用。2. 流式更新的效率对比更新时刻不使用 Tiling使用 Tiling片复用率T00:00计算 1 小时计算 12 个片0%首次T00:05重算 1 小时1000 事件计算 1 片 复用 11 片92%T00:10重算 1 小时1000 事件计算 1 片 复用 11 片92%T00:15重算 1 小时1000 事件计算 1 片 复用 11 片92%关键收益在流式会话期间所有片都驻留内存从而实现了超过 90% 的片级复用——这正是 Tiling 相对全量重算获得加速的根本原因。四、Tiling 算法Sawtooth锯齿窗窗口分片Feast 的 Tiling 实现名为Sawtooth Window Tiling核心分五步把事件按 hop 大小分区例如 5 分钟一个区间从物化窗口起点开始为每个 hop 计算累积尾部聚合cumulative tail aggregations做片减法当前片 - 前一片得到窗口化聚合存储 IR保证整体聚合合并正确物化时把窗口化聚合写入在线存储Online Store。算法收益查询时高效窗口结果预计算完毕无需查询时重扫存储开销极小只存 hop 大小的片数学上对全部聚合类型正确整体聚合靠 IR 分量加减保证。源码视角orchestrator 与 tile_subtraction从源码看算法被拆成两个纯 pandas 模块且与具体计算引擎解耦sdk/python/feast/aggregation/tiling/orchestrator.py 中的apply_sawtooth_window_tiling()负责生成累积片。其关键步骤包括把时间戳换算成毫秒并计算_hop_interval下界包含式(timestamp_ms // hop_size_ms) * hop_size_ms按group_by_keys [_hop_interval]分组聚合为代数聚合直接聚合、为整体聚合计算sum/count/sum_sq等 IR 列为每个实体补齐完整的 hop 网格没有数据的片补 0再对 IR 列做cumsum得到累积尾部值写入_tile_start/_tile_end元数据并把代数聚合的最终值直接赋给特征列。sdk/python/feast/aggregation/tiling/tile_subtraction.py 中的convert_cumulative_to_windowed()执行核心数学运算windowed_agg_at_T cumulative_tile_at_T - cumulative_tile_at_(T - window_size)对整体聚合avg/std/var而言它先对每个 IR 分量做减法再用窗口化后的 IR 重新计算最终值如avg windowed_sum / windowed_count最后通过deduplicate_keep_latest()保留每个实体的最新时间戳并丢弃_tile_start、_tile_end、_hop_interval等内部列产出可直接写入在线存储的结果。注意减法只接受_tile_end精确匹配window_start的片以避免算出超出请求窗口大小的错误窗口。五、配置在 StreamFeatureView 中启用 TilingTiling 对更新频繁的流式场景收益最大Feast 推荐在StreamFeatureView上启用。推荐用法Kafka 流式示例from feast import StreamFeatureView, Aggregation from feast.data_source import PushSource, KafkaSource from datetime import timedelta # 以 Kafka 流式数据源为例 customer_features StreamFeatureView( namecustomer_transaction_features, entities[customer], sourceKafkaSource( nametransactions_stream, kafka_bootstrap_serverslocalhost:9092, topictransactions, timestamp_fieldevent_timestamp, batch_sourcefile_source, # 用于历史数据 ), aggregations[ Aggregation(columnamount, functionsum, time_windowtimedelta(hours1), namesum_amount_1h), Aggregation(columnamount, functionavg, time_windowtimedelta(hours1), nameavg_amount_1h), Aggregation(columnamount, functionstd, time_windowtimedelta(hours1), namestd_amount_1h), ], timestamp_fieldevent_timestamp, onlineTrue, # Tiling 配置 enable_tilingTrue, # 流式场景的加速开关 tiling_hop_sizetimedelta(minutes5), # 更新频率 )何时启用流式数据源Kafka、Kinesis、PushSource高频更新每隔几分钟实时特征服务高吞吐事件处理。关键参数详解aggregations要计算的时间窗口聚合列表。每个Aggregation接受column要聚合的源列function聚合函数sum、avg、mean、min、max、count、stdtime_window聚合窗口的时长slide_intervalhop/滑动大小默认等于time_window——在 sdk/python/feast/aggregation/init.py 中未指定时self.slide_interval self.time_windowname可选输出特征名默认{function}_{column}如sum_amount设置自定义名称如namesum_amount_1h便于下游引用。timestamp_field时间戳列名指定aggregations时为必填否则 stream_feature_view.py 会抛出ValueError。enable_tiling是否启用 Tiling 优化默认False流式场景建议设为True。tiling_hop_size两个片之间的时间间隔默认 5 分钟。参数取舍更小 片更细处理窗口期内存可能更高更大 片更粗处理窗口期内存可能更低。一个必须遵守的约束hop_size 必须小于最小窗口从 sdk/python/feast/stream_feature_view.py 可以看到一个硬性校验当启用 Tiling 且存在time_window时Feast 会取所有聚合窗口的最小值若tiling_hop_size min(time_window)将直接抛出ValueError提示 If hop_size window_size, the tiling algorithm will produce incorrect results。因此配置时务必保证 hop 严格小于所有聚合窗口。计算引擎要求计算引擎支持情况Spark Compute Engine流式与批式完整支持Ray Compute Engine流式与批式完整支持Local Compute Engine不支持时间窗口聚合引擎侧接线pandas 与引擎转换Tiling 的实现刻意保持纯 pandas 内核 引擎薄转换层的架构┌─────────────────┐ │ Engine DataFrame│ (Spark/Ray 等) └────────┬────────┘ │ .toPandas() / .to_pandas() ▼ ┌─────────────────┐ │ Pandas DataFrame│ └────────┬────────┘ │ orchestrator.apply_sawtooth_window_tiling() ▼ ┌─────────────────┐ │ Cumulative │ (pandas含 _tile_start、_tile_end 与 IR 列) │ Tiles │ └────────┬────────┘ │ tile_subtraction.convert_cumulative_to_windowed() ▼ ┌─────────────────┐ │ Windowed │ (pandas含最终聚合) │ Aggregations │ └────────┬────────┘ │ spark.createDataFrame() / ray.from_pandas() ▼ ┌─────────────────┐ │ Engine DataFrame│ └─────────────────┘引擎侧的证据在 sdk/python/feast/infra/compute_engines/spark/nodes.py 中Spark 节点持有enable_tiling与hop_size参数当二者满足条件时导入apply_sawtooth_window_tiling与convert_cumulative_to_windowed并把默认 hop 设为timedelta(minutes5)Ray 计算引擎在 sdk/python/feast/infra/compute_engines/ray/nodes.py 中采用了同样的sawtooth分片逻辑。也就是说无论底层是 Spark 还是 Ray累积片生成与片减法都复用了同一份纯 pandas 实现。六、测试验证边界行为有据可依仓库单元测试 sdk/python/tests/unit/test_aggregation_ops.py 验证了一个关键边界对count_distinct调用 Tiling 元数据解析会抛出ValueError匹配信息 count_distinct does not support tiling印证了整体聚合中并非所有函数都能 Tiling的约束该测试还从feast.aggregation.tiling.base导入get_ir_metadata_for_aggregation说明 IR 分类逻辑是 Tiling 语义的核心入口。启用 Tiling 前建议在特征仓库中先跑通类似的最小化验证确认所选聚合函数均被支持。七、总结Tiling with Intermediate Representations 是 Feast 为流式时间窗口聚合提供的一套优化方案其要点可归纳为正确性来自 IRavg存sum countstd/var额外存sum_of_squares合并时做分量加减再复算最终值sum/count/max/min是代数聚合最终值即 IR加速来自复用Sawtooth 算法把窗口切成 hop 大小的片每次更新只计算 1 个新片、复用内存中其余 11 个片1 小时窗口 5 分钟 hop 场景下约 92% 复用率落地很简单在StreamFeatureView中设置enable_tilingTrue与tiling_hop_size默认 5 分钟并确保 hop 小于所有聚合窗口Spark 与 Ray 引擎均可使用Local 引擎不支持窗口聚合。对于从 Kafka/Kinesis/PushSource 高频更新、又对结果正确性有严格要求的实时特征场景Tiling 是一个又快又对的务实选择。【免费下载链接】feastThe Open Source Feature Store for AI/ML项目地址: https://gitcode.com/GitHub_Trending/fe/feast创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表