You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Snakemake流程中过滤保留每个ID的最新文件

解决Snakemake中按ID筛选最新日期文件的问题

你的报错核心原因是:Snakemake需要能从输出文件反向推导出输入中的通配符(比如ID),但直接用返回所有最新文件的函数作为输入时,无法建立这种绑定关系。下面给出两种可行的实现方案:

方案1:按ID拆分规则(推荐,适合需要单独处理每个ID的场景)

步骤1:提前收集所有唯一ID

在Snakefile开头先扫描目录,获取所有存在的ID,让Snakemake明确知道所有可能的通配符值:

import glob
import os

# 替换为你的输入文件路径
INPUT_DIR = "path/to/your/json/files"
all_ids = sorted({
    os.path.basename(f).split(".", 1)[0] 
    for f in glob.glob(f"{INPUT_DIR}/*.json")
})

步骤2:编写单个ID的最新文件筛选函数

针对每个ID,单独找到其对应的最新日期文件:

from datetime import datetime

def get_latest_file_for_id(wildcards):
    # 获取当前ID的所有文件
    id_files = glob.glob(f"{INPUT_DIR}/{wildcards.ID}.*.json")
    if not id_files:
        raise ValueError(f"No files found for ID: {wildcards.ID}")
    
    # 按日期排序(需根据你的日期格式调整strptime的参数,比如"%Y%m%d"或"%Y-%m-%d")
    sorted_files = sorted(
        id_files,
        key=lambda x: datetime.strptime(os.path.basename(x).split(".")[-2], "%Y-%m-%d")
    )
    # 返回最新的文件(排序后的最后一个)
    return sorted_files[-1]

步骤3:定义规则

创建一个汇总规则all,依赖所有ID的最新文件;再定义单个处理规则,绑定ID通配符:

rule all:
    input:
        expand("output/{ID}.latest.json", ID=all_ids)

rule process_latest_id:
    input:
        get_latest_file_for_id
    output:
        "output/{ID}.latest.json"
    shell:
        # 这里替换为你的实际处理命令,示例是复制最新文件到输出目录
        "cp {input} {output}"

方案2:用Checkpoint处理动态输入(适合需批量处理所有最新文件的场景)

如果不需要保留每个ID的单独输出,而是要将所有最新文件作为后续规则的输入,可以用Checkpoint先动态生成最新文件列表:

步骤1:定义Checkpoint收集最新文件

checkpoint collect_latest_files:
    output:
        touch("checkpoints/collect_latest.done")
    run:
        import os
        from datetime import datetime

        files = glob.glob(f"{INPUT_DIR}/*.json")
        id_to_latest = {}

        for f in files:
            fname = os.path.basename(f).replace(".json", "")
            id_str, date_str = fname.split(".", 1)
            file_date = datetime.strptime(date_str, "%Y-%m-%d")
            
            # 更新当前ID的最新文件
            if id_str not in id_to_latest or file_date > id_to_latest[id_str][0]:
                id_to_latest[id_str] = (file_date, f)
        
        # 将最新文件列表写入临时文件
        os.makedirs("checkpoints", exist_ok=True)
        with open("checkpoints/latest_files.txt", "w") as fp:
            fp.write("\n".join([f for _, f in id_to_latest.values()]))

步骤2:读取最新文件列表作为输入

def load_latest_files(wildcards):
    with open("checkpoints/latest_files.txt") as fp:
        return [line.strip() for line in fp if line.strip()]

rule process_all_latest:
    input:
        load_latest_files
    output:
        "output/all_latest_processed.json"
    shell:
        # 替换为你的批量处理命令,示例是合并所有文件
        "cat {input} > {output}"

rule all:
    input:
        "output/all_latest_processed.json"

关键说明

  • 两种方案都避免了Snakemake的通配符推导问题:方案1提前告知所有ID,让每个规则的输入输出都绑定明确的ID通配符;方案2用Checkpoint先完成动态筛选,再将静态列表传递给后续规则。
  • 注意调整日期解析的格式字符串(strptime的第二个参数),确保和你的文件名中的日期格式匹配。

内容的提问来源于stack exchange,提问作者Klumpi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 21:37:43