增量处理:快速管道更新
本文解释SEO管道如何利用增量处理在数秒内完成运行,而非数小时。
问题:完全重新处理速度缓慢
从头运行整个管道需要数小时:
-
步骤0(源数据嵌入):15分钟(65,000个产品)
-
步骤1(查询获取):10分钟(API调用)
-
步骤2(查询聚类):30分钟(65K×65K相似度计算)
-
步骤3(短语映射):20分钟(嵌入+匹配)
-
步骤4(产品匹配):45分钟(查询×产品)
-
步骤5(相关搜索):25分钟(查询×查询相似度)
总计:完整管道约2.5小时
问题:每日更新将浪费2.5小时重新计算未变化的数据。
解决方案:三层增量策略
我们采用三种技术来跳过不必要的工作:
1. 步骤跳过(粗粒度)
如果输出结果新鲜且脚本未更改,则跳过整个步骤。
2. 增量嵌入(中粒度)
仅嵌入新增/变更的条目,复用缓存的嵌入结果。
3. 检查点(细粒度)
在长时间操作中保存进度,失败时从检查点恢复。
步骤跳过策略
工作原理
每个步骤开始前检查:
输出是否存在? 若否,运行该步骤。
输出时效性:若超过7天,运行该步骤。
脚本是否变更? 若脚本在输出生成后被修改,运行该步骤。
所有检查通过? 跳过该步骤。
实现
def should_skip_step(output_path, script_path, days=7):
# 检查输出是否存在
if not os.path.exists(output_path):
# ...(实现细节省略)
使用方式
每个脚本在启动时检查:
from seo_common import should_skip_step
if should_skip_step(SEO_SOURCE_EMBEDDINGS_PATH, __file__):
print("✓ 跳过:输出结果新鲜且脚本未更改")
return
优势
快速每日运行:若数据未变,多数步骤可跳过
自动失效:脚本变更触发重新运行
可配置新鲜度:按步骤调整days参数
增量嵌入策略
工作原理
嵌入条目(产品、查询、短语)时:
加载缓存:读取先前嵌入的条目及其键值
比较键值:识别新增、变更和删除的条目
仅嵌入新增项:仅嵌入缓存中不存在的条目
合并:将缓存嵌入结果与新增嵌入结果按正确顺序合并
保存:写入更新后的缓存
实现
incremental_embed_with_keys函数处理此逻辑:
def incremental_embed_with_keys(
items, # 待嵌入的当前条目
keys, # 条目的唯一键值
# ...(实现细节省略)
缓存命中率
首次运行后的典型缓存命中率:
源数据(产品、部件、文章):
-
首次运行:0%(嵌入全部65,000个条目)
-
每日运行:99.5%(仅约300个新增/变更条目)
查询(来自GSC、广告、实时数据):
-
首次运行:0%(嵌入全部65,000个查询)
-
每日运行:99.2%(仅约500个新增查询)
短语映射:
-
首次运行:0%(嵌入全部5,000个短语)
-
每日运行:99.8%(仅约10个新增短语)
性能影响
首次运行(冷缓存):
-
源数据嵌入:15分钟(65,000个条目)
-
查询嵌入:10分钟(65,000个查询)
-
短语嵌入:2分钟(5,000个短语)
每日运行(热缓存):
-
源数据嵌入:10秒(300个条目,99.5%命中率)
-
查询嵌入:5秒(500个查询,99.2%命中率)
-
短语嵌入:1秒(10个短语,99.8%命中率)
加速效果:90-180倍提速
检查点策略
工作原理
针对长时间运行的操作(嵌入65,000个条目):
批处理:按批次处理条目(如每批1,000个)
保存检查点:每批完成后保存累计结果
失败恢复:若进程崩溃,从最后检查点恢复
最终保存:所有批次完成后保存完整结果
实现
检查点功能内置于incremental_embed_with_keys:
checkpoint_every = 1000 # 每1,000个条目保存一次
embeddings_list = []
# ...(实现细节省略)
优势
崩溃恢复:从最后检查点恢复而非重新开始
进度可见性:每1,000个条目可见进度
内存效率:批处理,无需一次性加载全部数据
管道集成应用
增量处理在多个步骤中使用:
步骤0:源数据嵌入
增量:仅嵌入新增/变更的产品、部件、文章
检查点:每1,000个条目保存一次
跳过逻辑:若输出<7天且脚本未更改则跳过
参见:源数据嵌入
步骤1:查询获取
增量:API调用仅获取新数据(自上次运行后)
跳过逻辑:若输出<1天则跳过
参见:查询获取
步骤3b:查询嵌入
增量:仅嵌入新查询
检查点:每1,000个查询保存一次
跳过逻辑:若输出<7天且脚本未更改则跳过
参见:查询嵌入
步骤4:短语映射扩展
增量:仅嵌入新短语
检查点:每1,000个短语保存一次
跳过逻辑:若输出<7天且脚本未更改则跳过
参见:短语到过滤器映射
步骤6:产品匹配
增量:仅匹配新查询
跳过逻辑:若输出<7天且脚本未更改则跳过
参见:产品匹配
配置
增量处理按步骤配置:
新鲜度阈值
# 若输出<7天则跳过(默认)
should_skip_step(output_path, script_path, days=7)
# 若输出<1天则跳过(针对频繁变化的数据)
should_skip_step(output_path, script_path, days=1)
检查点频率
# 每1,000个条目保存一次(默认)
incremental_embed_with_keys(..., checkpoint_every=1000)
# 每5,000个条目保存一次(处理更快,安全性降低)
incremental_embed_with_keys(..., checkpoint_every=5000)
批次大小
# 每批嵌入32个条目(默认,平衡型)
incremental_embed_with_keys(..., batch_size=32)
# 每批嵌入64个条目(GPU上更快,内存占用更高)
incremental_embed_with_keys(..., batch_size=64)
监控与调试
缓存统计
每个步骤打印缓存统计:
✓ 发现现有缓存,检查变更中...
现有:65,000个条目
当前:65,300个条目
复用:64,800个嵌入结果
新增:500个待嵌入条目
跳过消息
步骤被跳过时:
✓ 跳过0_embed_source_data.py:输出结果新鲜且脚本未更改。
检查点消息
长时间操作期间:
嵌入65,000个条目(每1,000个保存检查点)...
批次0-1000...
✓ 检查点已保存(累计1,000个)
批次1000-2000...
✓ 检查点已保存(累计2,000个)
...
参考文献
技术概念
相关文章
总结
增量处理使管道提速90-180倍:
三层策略:
-
✅ 步骤跳过(若输出新鲜则跳过整个步骤)
-
✅ 增量嵌入(仅嵌入新增/变更条目)
-
✅ 检查点(保存进度,失败时恢复)
性能表现:
-
✅ 首次运行:约2.5小时(完整管道)
-
✅ 每日运行:约5分钟(增量更新)
-
✅ 缓存命中率:首次运行后达99%以上
优势:
-
✅ 快速每日更新(分钟级而非小时级)
-
✅ 自动失效(脚本变更触发重新运行)
-
✅ 崩溃恢复(从检查点恢复)
-
✅ 内存高效(批处理)
该策略支持每日管道运行,无需为未变化数据浪费计算资源。