You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 14:06:02