Luigi数据管道中如何筛选新/修改文件并排除重命名文件?
问题解决:区分Luigi管道中的新文件/修改文件,以及工具选型疑问
一、修改代码区分新文件与重命名文件
原代码未检查哈希是否存在于数据库,导致把重命名文件(哈希已存在)和真正的新文件(哈希不存在)混为一谈。调整思路如下:
- 先从数据库提取所有已存哈希值存入集合,实现O(1)快速查找
- 遍历文件时,针对文件名不存在的情况,额外判断哈希是否在集合内:
- 哈希不在集合 → 标记为新文件
- 哈希在集合 → 判定为重命名文件,直接跳过
修改后的run方法代码:
def run(self): with ThreadPoolExecutor(max_workers=20) as executor: file_hashes = list(executor.map(self._calculate_file_hash, os.listdir(self.local_folder))) df = pd.read_sql_table('filesy_state', con=DB_ENGINES[self.engine_name]) # 提取数据库中已存在的哈希和文件名,存入集合优化查找效率 existing_hashes = set(df['file_hash'].tolist()) existing_filenames = set(df['file_name'].tolist()) new_files = [] modified_files = [] for file, file_hash in file_hashes: if file in existing_filenames: # 文件名存在,对比哈希判断是否修改 stored_hash = df[df['file_name'] == file]['file_hash'].iloc[0] if stored_hash != file_hash: modified_files.append({'file_name': file}) else: # 文件名不存在,检查哈希是否已存在 if file_hash not in existing_hashes: new_files.append({'file_name': file}) # 哈希存在则是重命名文件,跳过不记录 compared_files = {'modified_files': modified_files, 'new_files': new_files} with self.output().open('w') as outfile: json.dump(compared_files, outfile)
二、关于Luigi+SQLite是否过度复杂的疑问
不算过度复杂,反而很适合新手入门:
- SQLite:轻量无需单独部署,文件型数据库足够支撑小到中等规模的文件状态存储,上手成本极低,适合练手和小型项目。
- Luigi:如果只是简单的定时扫文件+解析写入,用crontab加Python脚本确实更直接;但如果未来有扩展需求(比如步骤依赖管理、失败重试、多任务并行、状态监控),Luigi的任务调度、依赖管理特性会大幅降低维护成本。现在用Luigi搭建基础框架,能提前熟悉数据管道核心概念,后续扩展更顺畅。
内容的提问来源于stack exchange,提问作者Bennet Turner
相关产品推荐
相关产品推荐

