Intake分区Parquet数据源如何按指定日期及日期范围加载数据?
前置说明
Intake的Parquet数据源默认采用懒加载机制,不会主动扫描路径下的所有文件元数据,因此未执行
discover()前npartitions返回0属于正常现象。
需求实现方案
两种实现路径,可根据场景选择:
方案1:直接匹配路径加载(无需全量扫描元数据)
直接通过urlpath规则匹配目标文件,不需要扫描全量1726个文件的元数据,加载效率更高,同时兼容平层/分层两种路径格式。
加载指定单日数据(如20180101)
复用现有catalog配置的写法:
import intake cat = intake.open_catalog('mycat.yaml') # **通配符兼容任意层级的父目录 source = cat.source_paritioned(urlpath="/some/data/**/20180101.parquet") # 加载到Pandas DataFrame df = source.read() # 也可转Dask DataFrame做后续处理 # ddf = source.to_dask()
加载指定日期范围数据(如20170601到20190223)
先生成日期范围内的所有日期字符串,再批量匹配路径:
import intake from datetime import datetime, timedelta cat = intake.open_catalog('mycat.yaml') start_dt = datetime(2017,6,1) end_dt = datetime(2019,2,23) # 生成范围内所有日期的字符串 date_list = [] for i in range((end_dt - start_dt).days + 1): current_dt = start_dt + timedelta(days=i) date_list.append(current_dt.strftime("%Y%m%d")) # 生成所有目标文件的匹配路径 target_paths = [f"/some/data/**/{dt}.parquet" for dt in date_list] source = cat.source_paritioned(urlpath=target_paths) df = source.read()
方案2:转Dask DataFrame后过滤(适合已有日期字段的场景)
如果你的Parquet文件内已经存储了日期字段,或者需要先加载全量分区元数据做多次过滤,可以使用to_dask()方法转换后做过滤:
import intake cat = intake.open_catalog('mycat.yaml') source = cat.source_paritioned() source.discover() ddf = source.to_dask() # 假设文件内日期字段名为dt,格式为'YYYYMMDD'字符串或日期类型 # 1. 过滤单日数据 single_day_df = ddf[ddf['dt'] == '20180101'].compute() # 2. 过滤日期范围数据 range_df = ddf[(ddf['dt'] >= '20170601') & (ddf['dt'] <= '20190223')].compute()
如果Parquet文件内没有存储日期字段,可以从文件路径中提取日期作为过滤字段:
import dask.dataframe as dd # 从文件路径中提取8位日期字符串 ddf = ddf.assign( dt=dd.str.slice(dd.str.split(ddf['__filename__'], '/').str[-1], 0, 8) )
内容的提问来源于stack exchange,提问作者Mikhail Shevelev
相关产品推荐
相关产品推荐

