如何在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
相关产品推荐
相关产品推荐

