增量处理:快速管道更新

本文解释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%以上

优势:

  • ✅ 快速每日更新(分钟级而非小时级)

  • ✅ 自动失效(脚本变更触发重新运行)

  • ✅ 崩溃恢复(从检查点恢复)

  • ✅ 内存高效(批处理)

该策略支持每日管道运行,无需为未变化数据浪费计算资源。


← 返回文档索引