jose-torres commented on code in PR #57365:
URL: https://github.com/apache/spark/pull/57365#discussion_r3626804067
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala:
##########
@@ -811,6 +811,93 @@ case class Scd2BatchProcessor(
staged.select(outputColumns.toImmutableArraySeq: _*)
}
+ /**
+ * Drop delete-encoded rows (tombstones and decomposition tails) that became
redundant after
+ * reconciliation.
+ *
+ * A tombstone or decomposition tail is redundant when the immediately
preceding row's reconciled
+ * [[endAtColName]] equals its own sequence: the preceding upsert already
encodes the delete
+ * boundary, so the standalone delete-encoded row should no longer be routed
to aux.
+ *
+ * For example, an open upsert `[startAt=10, endAt=null)` followed by a
tombstone at `15`
+ * reconciles into a closed upsert `[10, 15)`, making the tombstone
redundant. Likewise, an
+ * existing closed `[10, 20)` bisected by an event at `15` reconciles the
event into `[15, 20)`,
+ * making the `[null, 20)` decomposition tail redundant.
+ */
+ 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)
+ )
+
+ withWindowCols
+ .filter(!F.col(Scd2BatchProcessor.isRedundantDeleteEncodingColName))
+ .drop(
+ Scd2BatchProcessor.previousEndAtColName,
+ Scd2BatchProcessor.isRedundantDeleteEncodingColName
+ )
+ }
+
+ /**
+ * Convert surviving decomposition tails into tombstones.
Review Comment:
Another thing to handle in a post-hoc pass after implementation: we should
make sure that these two are actually handled in a different enough way to
require two different representations. It kinda sounds to me like a
decomposition tail is just the within-batch version of a tombstone.
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala:
##########
@@ -811,6 +811,93 @@ case class Scd2BatchProcessor(
staged.select(outputColumns.toImmutableArraySeq: _*)
}
+ /**
+ * Drop delete-encoded rows (tombstones and decomposition tails) that became
redundant after
+ * reconciliation.
+ *
+ * A tombstone or decomposition tail is redundant when the immediately
preceding row's reconciled
+ * [[endAtColName]] equals its own sequence: the preceding upsert already
encodes the delete
+ * boundary, so the standalone delete-encoded row should no longer be routed
to aux.
+ *
+ * For example, an open upsert `[startAt=10, endAt=null)` followed by a
tombstone at `15`
Review Comment:
The way this is stated is confusing to me. Why does the history record
identity follow the prefix of the range when closing an active row but the
suffix when closing an inactive row?
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala:
##########
@@ -811,6 +811,93 @@ case class Scd2BatchProcessor(
staged.select(outputColumns.toImmutableArraySeq: _*)
}
+ /**
+ * Drop delete-encoded rows (tombstones and decomposition tails) that became
redundant after
+ * reconciliation.
+ *
+ * A tombstone or decomposition tail is redundant when the immediately
preceding row's reconciled
+ * [[endAtColName]] equals its own sequence: the preceding upsert already
encodes the delete
+ * boundary, so the standalone delete-encoded row should no longer be routed
to aux.
+ *
+ * For example, an open upsert `[startAt=10, endAt=null)` followed by a
tombstone at `15`
+ * reconciles into a closed upsert `[10, 15)`, making the tombstone
redundant. Likewise, an
+ * existing closed `[10, 20)` bisected by an event at `15` reconciles the
event into `[15, 20)`,
+ * making the `[null, 20)` decomposition tail redundant.
+ */
+ 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)
+ )
+
+ withWindowCols
Review Comment:
I suppose the implementation of this method doesn't end up depending on the
row identity confusion above, but I still think it's worth clearing up.
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala:
##########
@@ -811,6 +811,93 @@ case class Scd2BatchProcessor(
staged.select(outputColumns.toImmutableArraySeq: _*)
}
+ /**
+ * Drop delete-encoded rows (tombstones and decomposition tails) that became
redundant after
+ * reconciliation.
+ *
+ * A tombstone or decomposition tail is redundant when the immediately
preceding row's reconciled
+ * [[endAtColName]] equals its own sequence: the preceding upsert already
encodes the delete
+ * boundary, so the standalone delete-encoded row should no longer be routed
to aux.
+ *
+ * For example, an open upsert `[startAt=10, endAt=null)` followed by a
tombstone at `15`
+ * reconciles into a closed upsert `[10, 15)`, making the tombstone
redundant. Likewise, an
+ * existing closed `[10, 20)` bisected by an event at `15` reconciles the
event into `[15, 20)`,
+ * making the `[null, 20)` decomposition tail redundant.
+ */
+ 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)
+ )
+
+ withWindowCols
+ .filter(!F.col(Scd2BatchProcessor.isRedundantDeleteEncodingColName))
+ .drop(
+ Scd2BatchProcessor.previousEndAtColName,
+ Scd2BatchProcessor.isRedundantDeleteEncodingColName
+ )
+ }
+
+ /**
+ * Convert surviving decomposition tails into tombstones.
+ *
+ * A decomposition tail that survives deletion cleanup is an unmatched
delete boundary; setting
+ * [[recordStartAtFieldName]], and [[startAtColName]] to the tail's end
sequence lets downstream
+ * aux handling preserve it as a tombstone.
+ *
+ * For example, if an existing closed upsert `[startAt=10, endAt=20)` is
bisected by an event at
+ * `15`, decomposition first produces an open head `[10, null)` plus a tail
`[null, 20)`;
+ * reconciliation may close the head as `[10, 15)`, leaving the tail's
boundary at `20` to be
+ * promoted into a tombstone.
+ */
+ private[autocdc] def promoteDecompositionTailsToTombstones(
+ reconciledDf: DataFrame): DataFrame = {
+ val recordStartAt =
+
Scd2BatchProcessor.recordStartAtOf(F.col(AutoCdcReservedNames.cdcMetadataColName))
+ val startAt = F.col(Scd2BatchProcessor.startAtColName)
+ val endAt = F.col(Scd2BatchProcessor.endAtColName)
+ val isDecompositionTail =
+ RowClassifier.isDecompositionTail(Scd2IntervalColumns(recordStartAt,
startAt, endAt))
+
+ val outputColumns = reconciledDf.columns.map {
+ case c if c == AutoCdcReservedNames.cdcMetadataColName =>
+ val metadata = reconciledDf.schema(c).metadata
+ F.when(
+ isDecompositionTail,
+ Scd2BatchProcessor.constructCdcMetadataCol(
+ recordStartAt = endAt,
+ sequencingType = resolvedSequencingType
+ )
+ ).otherwise(F.col(c)).as(c, metadata)
+ case c if c == Scd2BatchProcessor.startAtColName =>
+ val metadata = reconciledDf.schema(c).metadata
+ F.when(isDecompositionTail, endAt).otherwise(startAt).as(c, metadata)
+ case c =>
+ F.col(QuotingUtils.quoteIdentifier(c))
Review Comment:
I guess we only need the quote here because we know the exact names of the
other cases and they don't contain special characters?
--
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]