【问题标题】:PySpark : Optimize read/load from Delta using selected columns or partitionsPySpark:使用选定的列或分区优化从 Delta 读取/加载
【发布时间】: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,给定:

  1. 只需要特定的日期范围
  2. 只需要列的子集

我尝试过的当前方法是:

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


【解决方案1】:

在您的情况下,不需要额外的步骤。 Spark 将负责优化。由于当您尝试使用分区列 dt 作为过滤条件查询数据集时,您已经根据列 dt 对数据集进行了分区。 Spark 仅从源数据集中加载与过滤条件匹配的数据子集,在您的情况下为 dt > '2020-06-20'

Spark 在内部进行基于优化的分区修剪。

【讨论】:

    【解决方案2】:

    无需 SQL 即可完成此操作..

    from pyspark.sql import functions as F
    
    df = spark.read.format("delta").load(path_to_data).filter(F.col("dt_col") > F.lit('2020-06-20'))
    

    虽然对于这个示例,您可能需要做一些比较日期的工作。

    【讨论】:

      猜你喜欢
      • 2018-11-14
      • 2021-12-05
      • 1970-01-01
      • 1970-01-01
      • 2021-01-05
      • 2020-07-23
      • 1970-01-01
      • 1970-01-01
      • 2023-03-02
      相关资源
      最近更新 更多