[
https://issues.apache.org/jira/browse/SPARK-58979?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated SPARK-58979:
-----------------------------------
Labels: pull-request-available structured-streaming (was:
structured-streaming)
> ColumnPruning disables late-event filtering when event time is dropped after
> dropDuplicates
> -------------------------------------------------------------------------------------------
>
> Key: SPARK-58979
> URL: https://issues.apache.org/jira/browse/SPARK-58979
> Project: Spark
> Issue Type: Bug
> Components: Optimizer, Structured Streaming
> Affects Versions: 5.0.0
> Reporter: Jonathan Gao
> Priority: Major
> Labels: pull-request-available, structured-streaming
>
> Consider the following streaming queries with identical input and watermark
> configuration.
> PRUNE-1 projects the event-time column away after dropDuplicates:
> {code:scala}
> val input = MemoryStream[(String, Long)]
> val result = input.toDF()
> .selectExpr("_1 AS id", "CAST(_2 AS TIMESTAMP) AS ts")
> .withWatermark("ts", "10 seconds")
> .dropDuplicates("id")
> .select("id")
> testStream(result, Append)(
> AddData(input, ("a", 1000000L)),
> CheckNewAnswer("a"),
> AddData(input, ("b", 1000L)),
> CheckNewAnswer())
> {code}
> PRUNE-2 retains the event-time column above dropDuplicates:
> {code:scala}
> val input = MemoryStream[(String, Long)]
> val result = input.toDF()
> .selectExpr("_1 AS id", "CAST(_2 AS TIMESTAMP) AS ts")
> .withWatermark("ts", "10 seconds")
> .dropDuplicates("id")
> .selectExpr("id", "CAST(ts AS LONG) AS tsl")
> testStream(result, Append)(
> AddData(input, ("a", 1000000L)),
> CheckNewAnswer(("a", 1000000L)),
> AddData(input, ("b", 1000L)),
> CheckNewAnswer())
> {code}
> PRUNE-2 correctly drops the late row. PRUNE-1 currently emits id "b". The
> final projection therefore changes the late-event behavior of the same
> streaming deduplication operation.
> The root cause is that Deduplicate.references contains only the deduplication
> keys. ColumnPruning consequently inserts a Project below Deduplicate that
> removes the watermark-tagged event-time attribute. StreamingDeduplicateExec
> searches its child output for that attribute to construct the late-event
> predicate; after pruning, no predicate is created and the late row is emitted.
> The proposed fix is to include watermark-tagged child attributes in
> references for streaming Deduplicate, preventing ColumnPruning from removing
> an implicit input dependency. This follows the established fix for
> DeduplicateWithinWatermark in SPARK-50492. Batch Deduplicate continues to
> reference only its keys.
> Related: https://issues.apache.org/jira/browse/SPARK-50492
> Proposed patch branch:
> https://github.com/nahtonaj/spark/tree/fix/spark-dedup-watermark-pruning
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]