如何在Amundsen中并行化Snowflake元数据摄入任务?遇文件缺失错误
解决Amundsen并行Snowflake元数据摄入时的FileNotFoundError问题
看起来你在并行处理多个Snowflake账号的元数据摄入到Neo4j时,遇到了FileNotFoundError(找不到临时CSV文件),还偶尔会跳过部分数据库的情况。这本质是多进程共享同一个临时目录引发的资源竞争问题,我来帮你拆解原因并给出具体解决方案:
问题根源
Amundsen的neo4j_csv_publisher会在默认临时目录(比如/var/tmp/amundsen/)生成节点/关系的CSV文件,当多个进程同时读写这个目录时,会出现两种核心问题:
- 不同进程生成的CSV文件名冲突,后生成的文件覆盖了先生成的,导致读取时找不到目标文件
- 某个进程完成任务后自动清理临时文件,误删了其他进程还在使用的文件
解决方案
1. 为每个进程分配独立的临时目录(最推荐)
让每个Snowflake账号的处理进程使用专属的临时目录,彻底避免文件冲突。修改你的multiprocessing_snowflake_accounts函数,为每个账号创建独立的临时路径,并配置给Amundsen:
def multiprocessing_snowflake_accounts(ac_key, ac_config): import tempfile import shutil from pyhocon import ConfigFactory # 为当前账号创建唯一的临时目录 temp_dir = tempfile.mkdtemp(prefix=f"amundsen_sf_{ac_key}_") # 覆盖Amundsen的临时目录配置 custom_config = ConfigFactory.from_dict({ 'publisher.neo4j_csv_publisher.temp_file_dir': temp_dir, # 保留其他原有配置项... }) # 这里执行你的元数据摄入逻辑,使用custom_config替代原有配置 # 比如初始化extractor、transformer、publisher时传入这个配置 # ... 你的现有业务代码 ... # 任务完成后清理临时目录(可选,系统重启后也会自动清理) shutil.rmtree(temp_dir)
2. 限制并发进程数,避免资源过载
如果账号数量较多,同时启动所有进程会导致系统资源耗尽,间接引发文件操作异常。改用multiprocessing.Pool来控制并发数,比如限制为CPU核心数的一半:
def process_all_snowflake_accounts(): snowflake_config = read_snowflake_configuration() start_time = time.time() # 限制并发数,根据服务器资源调整 max_processes = multiprocessing.cpu_count() // 2 with multiprocessing.Pool(processes=max_processes) as pool: # 批量执行每个账号的处理任务 pool.starmap(multiprocessing_snowflake_accounts, snowflake_config.items()) print("CPU Unit: ", multiprocessing.cpu_count()) print('****************************************************************') print('Total time taken: ', time.time() - start_time) print('****************************************************************')
3. 禁用自动清理临时文件(临时调试用)
如果需要排查问题,可以先禁用Amundsen的临时文件自动清理功能,看看文件是否真的生成了。在配置里添加:
'publisher.neo4j_csv_publisher.keep_files': True
不过这个不适合生产环境,只是用来验证问题根源。
额外说明
你遇到的“偶尔跳过部分账号数据库”的问题,大概率是因为进程在遇到FileNotFoundError后崩溃,导致该账号后续的数据库处理被中断。解决了文件冲突问题后,这个现象也会随之消失。
内容的提问来源于stack exchange,提问作者Nishant Sharma
相关产品推荐
相关产品推荐

