【发布时间】:2020-06-23 17:13:01
【问题描述】:
我正在尝试将数据从 Delta 加载到 pyspark 数据帧中。
path_to_data = 's3://mybucket/daily_data/'
df = spark.read.format("delta").load(path_to_data)
现在基础数据按日期划分为
s3://mybucket/daily_data/
dt=2020-06-12
dt=2020-06-13
...
dt=2020-06-22
有没有办法将读取优化为 Dataframe,给定:
- 只需要特定的日期范围
- 只需要列的子集
我尝试过的当前方法是:
df.registerTempTable("my_table")
new_df = spark.sql("select col1,col2 from my_table where dt_col > '2020-06-20' ")
# dt_col is column in dataframe of timestamp dtype.
在上述状态下,Spark 是否需要加载整个数据,根据日期范围过滤数据,然后过滤所需的列?是否可以在 pyspark 读取中进行任何优化,以加载数据,因为它已经分区了?
某事在线:
df = spark.read.format("delta").load(path_to_data,cols_to_read=['col1','col2'])
or
df = spark.read.format("delta").load(path_to_data,partitions=[...])
【问题讨论】:
-
如果您在查询中调用
.explain,则生成的查询计划应提及引用您的dt列的PartitionFilter。您可以将此查询计划添加到您的帖子中吗?
标签: python apache-spark pyspark delta-lake