szehon-ho commented on code in PR #57365:
URL: https://github.com/apache/spark/pull/57365#discussion_r3632403056
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala:
##########
@@ -811,6 +811,107 @@ case class Scd2BatchProcessor(
staged.select(outputColumns.toImmutableArraySeq: _*)
}
+ /**
+ * Drop delete-encoded rows (tombstones and decomposition tails) that became
redundant after
+ * reconciliation.
+ *
+ * The rule is the same for both row kinds: a delete-encoded row encodes a
delete boundary in its
+ * [[endAtColName]], and it is redundant once the immediately preceding
reconciled upsert closes
+ * (its [[endAtColName]]) exactly on that same boundary. When they coincide,
the closed upsert
+ * already carries the boundary, so the standalone delete-encoded row can be
dropped rather than
+ * routed to aux. Only the preceding upsert's `endAt` matters here; its
`startAt` (identity) is
+ * irrelevant to the comparison.
+ *
+ * Both examples reduce to that single check - the preceding upsert's
reconciled `endAt` equals
+ * the delete-encoded row's boundary:
+ * - Tombstone: an open upsert `[startAt=10, endAt=null)` followed by a
tombstone at `15`
+ * reconciles into a closed upsert `[10, 15)`. Its `endAt=15` matches
the tombstone's
+ * boundary `15`, so the tombstone is redundant.
+ * - Decomposition tail: an existing closed `[10, 20)` is bisected by an
out-of-order event at
+ * `15`. Decomposition first splits the original row into an open head
`[10, null)` plus a
+ * decomposition tail `[null, 20)` - a synthetic row with a null
recordStartAt that carries
+ * the original right boundary `20` so it can be re-closed later. The
bisecting event then
+ * reconciles into `[15, 20)`, whose `endAt=20` matches the tail's
boundary `20`, so the tail
+ * is redundant. (A tail that does *not* coincide with a preceding
upsert survives and is
+ * later turned into a tombstone by
[[promoteDecompositionTailsToTombstones]].)
+ */
+ private[autocdc] def dropLeftoverDeletesPostReconciliation(
+ reconciledDf: DataFrame): DataFrame = {
+ val recordStartAt =
+
Scd2BatchProcessor.recordStartAtOf(F.col(AutoCdcReservedNames.cdcMetadataColName))
+ val startAt = F.col(Scd2BatchProcessor.startAtColName)
+ val endAt = F.col(Scd2BatchProcessor.endAtColName)
+
+ // Both tombstones and decomposition tails encode a delete boundary in
their `endAt`. Either
+ // becomes redundant when the immediately preceding upsert was reconciled
to close exactly on
+ // that boundary, since the resulting closed upsert already carries it.
+ val row = Scd2IntervalColumns(recordStartAt, startAt, endAt)
+ val isTombstone = RowClassifier.isTombstone(row)
+ val isDecompositionTail = RowClassifier.isDecompositionTail(row)
+ val isDeleteEncodedRow = isTombstone || isDecompositionTail
+
+ val withWindowCols = reconciledDf
+ .withColumn(
+ Scd2BatchProcessor.previousEndAtColName,
+ F.lag(endAt, 1).over(orderChronologicallyPerKeyWindow)
+ )
+ .withColumn(
+ Scd2BatchProcessor.isRedundantDeleteEncodingColName,
+ isDeleteEncodedRow && (F.col(Scd2BatchProcessor.previousEndAtColName)
<=> endAt)
Review Comment:
Could we defensively verify that the preceding row is an upsert? Adjacent
delete-encoded rows with the same boundary should not occur because upstream
cleanup prevents them, but checking here would make this method enforce its
documented invariant directly. For example:
```scala
val previous = row.lagBy(1, orderChronologicallyPerKeyWindow)
val isRedundantDeleteEncoding =
isDeleteEncodedRow &&
// Defensive: the predecessor should always be an upsert because
upstream cleanup prevents
// adjacent delete-encoded rows with the same boundary.
RowClassifier.isUpsertRepresentingRow(previous) &&
(previous.endAt <=> endAt)
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]