如何在Kedro中集成网络下载的tar.gz格式TSV数据集?
最优方案:Kedro集成lastfm-1K的tar.gz数据集
1. 配置HTTP数据集下载tar.gz
在conf/base/catalog.yml中添加用于下载的HTTP数据集,直接指向tar.gz文件的下载地址:
lastfm_raw_tar: type: kedro_datasets.http.HTTPDataSet url: "http://ocelma.net/MusicRecommendationDataset/lastfm-1K/userid-timestamp-artid-artname-traid-traname.tsv.tar.gz" load_args: timeout: 60 # 按需调整超时时间
2. 编写节点实现解压与TSV提取
在src/<你的项目名>/nodes/data_processing.py中编写解压函数,用Python原生模块处理tar.gz文件:
import tarfile from pathlib import Path def extract_tsv_from_tar(tar_path: str, output_dir: str) -> str: """从tar.gz中提取TSV文件并返回文件路径""" output_path = Path(output_dir) output_path.mkdir(parents=True, exist_ok=True) with tarfile.open(tar_path, "r:gz") as tar: # 筛选tar包中的TSV文件 tsv_files = [member for member in tar.getmembers() if member.name.endswith(".tsv")] if not tsv_files: raise ValueError("tar.gz文件中未找到TSV格式文件") tar.extract(tsv_files[0], path=output_dir) return str(output_path / tsv_files[0].name)
接着在src/<你的项目名>/pipeline.py中注册该节点:
from kedro.pipeline import Pipeline, node from .nodes.data_processing import extract_tsv_from_tar def create_pipeline(**kwargs) -> Pipeline: return Pipeline( [ node( func=extract_tsv_from_tar, inputs=["lastfm_raw_tar", "params:tsv_extraction_output_dir"], outputs="lastfm_extracted_tsv_path", name="extract_lastfm_tsv_node", ), ] )
同时在conf/base/parameters.yml中配置输出目录参数:
tsv_extraction_output_dir: "data/01_raw/lastfm_extracted"
3. 配置TSV数据集加载提取后的文件
利用Kedro的CSVDataSet(支持指定分隔符)加载TSV文件,在conf/base/catalog.yml中添加:
lastfm_cleaned_data: type: kedro_datasets.csv.CSVDataSet filepath: "{{ lastfm_extracted_tsv_path }}" load_args: sep: "\t" header: None # 原lastfm数据集无表头,按需调整 names: ["user_id", "timestamp", "artist_id", "artist_name", "track_id", "track_name"] # 手动指定列名 save_args: index: False
4. 运行验证
执行kedro run命令,Kedro会自动完成下载tar.gz、解压提取TSV、加载数据集的全流程,后续可直接在其他节点中使用lastfm_cleaned_data进行分析。
关键优化点
- 缓存配置:给HTTPDataSet添加缓存参数,避免重复下载大文件:
lastfm_raw_tar: # 其他配置不变 cache_args: cache_folder: "data/00_raw_cache" cache_key: "lastfm_tar"
- 异常处理:在解压节点中添加文件损坏、格式不符等异常捕获逻辑,提升鲁棒性;
- 资源控制:若数据集较大,可在节点中添加内存占用监控,避免OOM问题。
内容的提问来源于stack exchange,提问作者gaut
相关产品推荐
相关产品推荐

