【问题标题】:Pushdown filter in case of spark structured Delta streaming火花结构化 Delta 流的情况下的下推滤波器
【发布时间】:2021-05-26 09:33:43
【问题描述】:

我有一个用例,我们需要将开源增量表流式传输到多个查询中,并在其中一个分区列上进行过滤。 例如,。 给定按年份列分区的 Delta 表。

Streaming query 1
spark.readStream.format("delta").load("/tmp/delta-table/").
where("year= 2013")

Streaming query 2
spark.readStream.format("delta").load("/tmp/delta-table/").
where("year= 2014")

实物图在流后显示过滤器。

> == Physical Plan == Filter (isnotnull(year#431) AND (year#431 = 2013))
> +- StreamingRelation delta, []

我的问题是下推谓词是否适用于 Delta 中的流式查询? 我们可以仅从 Delta 流式传输特定分区吗?

【问题讨论】:

    标签: apache-spark delta-lake


    【解决方案1】:

    如果列已分区,则仅扫描所需的分区。

    让我们创建分区和非分区增量表并执行结构化流。

    分区增量表流式传输:

    val spark = SparkSession.builder().master("local[*]").getOrCreate()
    spark.sparkContext.setLogLevel("ERROR")
    import spark.implicits._
        
    //sample dataframe
    val df = Seq((1,2020),(2,2021),(3,2020),(4,2020),
    (5,2020),(6,2020),(7,2019),(8,2019),(9,2018),(10,2020)).toDF("id","year")
        
    //partionBy year column and save as delta table
    df.write.format("delta").partitionBy("year").save("delta-stream")
        
    //streaming delta table
    spark.readStream.format("delta").load("delta-stream")
    .where('year===2020)
    .writeStream.format("console").start().awaitTermination()
    

    上述流式查询的物理方案:注意partitionFilters

    非分区增量表流式传输:

    df.write.format("delta").save("delta-stream")
    
    spark.readStream.format("delta").load("delta-stream")
        .where('year===2020)
        .writeStream.format("console").start().awaitTermination()
    

    上述流式查询的物理方案:注意pushFilters

    【讨论】:

    • 您使用的是开源版本还是 Databricks 版本?在 OSS 中,下推过滤器不存在。更新问题以提及开源版本。
    • @AmitJoshi- 开源
    • 你能告诉我现在使用的 Delta core 的版本吗?
    • spark 3.0.1 和 delta-core 0.7.0
    • 我不确定,但我看不到推送的过滤器。您能否将代码粘贴到打印执行计划的位置。可能对我来说它被截断了。
    猜你喜欢
    • 2021-05-31
    • 2019-12-14
    • 2017-08-04
    • 2018-07-12
    • 1970-01-01
    • 2020-02-12
    • 2020-08-18
    • 2019-06-08
    • 2015-12-08
    相关资源
    最近更新 更多