Skip to content

[SPARK-57222][SDP] Implement SCD2 Batch Processor; Decompose affected rows - #56311

Closed
AnishMahto wants to merge 42 commits into
apache:masterfrom
AnishMahto:SPARK-57222-SCD2-decompose-affected-rows
Closed

AnishMahto wants to merge 42 commits into
apache:masterfrom
AnishMahto:SPARK-57222-SCD2-decompose-affected-rows

Conversation

@AnishMahto

@AnishMahto AnishMahto commented Jun 3, 2026 •

Copy link
Copy Markdown
Contributor

Approved AutoCDC SPIP: https://lists.apache.org/thread/j6sj9wo9odgdpgzlxtvhoy7szs0jplf7


What changes were proposed in this pull request?

Preamble:

The SCD type 2 flow is a foreachBatch streaming query on an input change-data-feed, and is responsible for reconciling the incoming change data onto some target table that follows SCD2 replication semantics.

SCD2 flows also maintain an "auxiliary" table to keep track of early-arriving out-of-order received events state. Each microbatch will need to reconcile against this auxiliary table as well, and update the auxiliary table's state appropriately for future microbatches.

Decompose affected rows:

Given the set of affected rows in the current microbatch execution - incoming rows in the microbatch, affected rows from aux table, affected rows from target table - the first step in microbatch reconciliation is decomposing closed historical rows that are being bisected by the microbatch.

A closed historical row is a row in the target table that has a non-null start-at and end-at. It's possible an incoming upsert/delete in the microbatch lands with a sequence in between an existing closed row's start/end at (i.e is a late-arriving event), bisecting it.

Decomposing a closed row means exactly this - bisecting the closed interval into a left and right end point, called the decomposed head and tail of the original closed row respectively. The head represents some past upsert event, the tail represents some past delete or upsert-overtake event.

Once a closed row is decomposed into its end points, it can either coalesce with other endpoints/events from the full set of affected rows to form a new historical row, or it can be demoted back to the aux table as a tombstone or no-op upsert.

Assert well formed rows post-decomposition
Given the decomposed affected rows, assert each classifies into one of the known row-types; tombstone, synthetic decomposition tail, or upsert representing row (which can be further classified as a valid open or closed upsert).

Drop redundant rows post-decomposition:

Decomposition transforms a single row into two synthetic children rows (decomposition head and tail), and its possible the resulting decomposition tail is made logically redundant by the incoming microbatch.

It's also possible that there are duplicate events by key+sequencing, either both introduced in the same microbatch or across the incoming and a past microbatch. Post-decomposition is a good time to reconcile these duplicates, because by this point we have the maximal possible set of rows that will be merged back into the target/aux tables.

We must drop redundant rows prior to reconciling the new start/end-ats for this decomposed set of microbatch-affected rows, to:

  1. Prevent zero width rows (start at == end at) from ever appearing in the target table.
  2. Ensure reconciliation of a row's start/end at is fully derived from at most one other row in the decomposed set.

Why are the changes needed?

Needed for core SCD2 reconciliation logic.

Does this PR introduce any user-facing change?

No, SCD2 is unreleased.

How was this patch tested?

Unit tests added to Scd2BatchProcessorSuite.

Was this patch authored or co-authored using generative AI tooling?

Co-authored with Claude Opus 4.7

Comment on lines +543 to +601
/**
* Asserts that every row in `decomposedRowsPerKey` conforms to one of the four canonical
* post-decomposition shapes - tombstone, open upsert, closed upsert, or decomposition
* tail - and is otherwise a structural identity transform.
*
* @param decomposedRowsPerKey
* the output of [[decomposeOutOfOrderRows]]: a dataframe conforming to the canonical
* SCD2 row schema `[user_cols..., [[startAtColName]], [[endAtColName]],
* [[cdcMetadataColName]]]`.
* @return
* a dataframe with the exact same schema and rows as the input. Failing the
* well-formedness check is treated as an internal-invariant violation: at execution
* time, the first ill-formed row encountered aborts the query with a SparkRuntimeException
*/
private[autocdc] def assertWellFormedRowsPostDecomposition(
decomposedRowsPerKey: DataFrame,
batchId: Long
): DataFrame = {
val recordStartAtField =
Scd2BatchProcessor.recordStartAtOf(F.col(AutoCdcReservedNames.cdcMetadataColName))
val startAtCol = F.col(Scd2BatchProcessor.startAtColName)
val endAtCol = F.col(Scd2BatchProcessor.endAtColName)

val isWellFormedRow =
RowClassifier.isDecompositionTail(recordStartAtField, startAtCol, endAtCol) ||
RowClassifier.isTombstone(recordStartAtField, startAtCol, endAtCol) ||
RowClassifier.isUpsertRepresentingRow(recordStartAtField, startAtCol, endAtCol)

def stringOrNullLit(c: Column): Column = F.coalesce(c.cast(StringType), F.lit("null"))
val malformedRowDiagnostic = F.concat(
F.lit(
s"During SCD2 reconciliation of microbatch [id=${batchId}], encountered a " +
"post-decomposition row of unexpected shape:"
),
F.lit(s" ${Scd2BatchProcessor.recordStartAtFieldName}="),
stringOrNullLit(recordStartAtField),
F.lit(s", ${Scd2BatchProcessor.startAtColName}="),
stringOrNullLit(startAtCol),
F.lit(s", ${Scd2BatchProcessor.endAtColName}="),
stringOrNullLit(endAtCol),
F.lit(".")
)

val internalErrorOnMalformed = ExpressionUtils.column(
If(
predicate = ExpressionUtils.expression(isWellFormedRow),
trueValue = Literal(null, BooleanType),
falseValue = RaiseError(
Literal("INTERNAL_ERROR"),
CreateMap(Seq(
Literal("message"),
ExpressionUtils.expression(malformedRowDiagnostic)
)),
BooleanType
)
)
)
decomposedRowsPerKey.filter(internalErrorOnMalformed.isNull)
}

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Continuing discussion from the previous PR, here's the first internal validation logic that I'm considering adding to SCD2 microbatch reconciliation.

Internal validation here means we are asserting on properties that we expect to be true if our implementation is correct, and all of the invariants we expect to be true actually hold.

There's a real argument to be made for not having such validation. If the validation throws here, it indicates either the implementation is wrong OR the user has gone out of their way to mutate the aux/target in a way that is invalid with the invariants required by the algorithm. But we cannot distinguish between the two cases, so there isn't really any meaningful action users can take if this throws, other than trying a full-refresh and hoping the issue was transient.

There also is some runtime penalty for executing this validation, although I'd argue it's likely negligible. The filter operation here is lazy, so we shouldn't be incurring an additional full dataframe scan.

From an implementation hygiene perspective, there's also an argument that invariants should be asserted and contractualized in tests and not in execution. While I buy in to that logic, the SCD2 algorithm is sophisticated enough where we cannot reasonably expect test all meaningfully different input sets that stress-test the assumed invariants. At best we can have the random-data fuzz test, but that is still non-deterministic random sampling.

Given all of this, reasons why I would argue for keeping this runtime-validation:

  • Although the exception is non-actionable for users, it's an extra guard to prevent us from accidentally corrupting user data if the implementation is incorrect or invariants are violated. At best the algorithm would throw at a later point anyway, at worst we corrupt the target/aux with incorrect data silently
  • It's not super expensive
  • The validation logic is self-documenting in a meaningful way. We can write as many comments as we want to try to document implicit invariants, but readers will have no way of knowing whether the comments themselves are still accurate, and readers will have to go through the full reasoning process by reading through multiple documented invariants across this algorithm. In contrast, post-validation, a reader can now genuinely trust the asserted properties, since the implementation would have thrown otherwise before that point.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I agree on balance that this makes sense. This is a compromise being forced by the untyped nature of dataframes, in a more strongly typed system we would 1000% do what's necessary to make this check happen at compile time.

@AnishMahto

Copy link
Copy Markdown
Contributor Author

@jose-torres @hhhhhazelnut for review.

About ~400 LOC of actual logic, which is still fairly large - sorry in advance. I don't think I can reasonably break this up into further PRs without losing coherency, but hopefully the PR description makes it pretty clear what's going on.

The logic here isn't ground breaking but is getting into the heart of SCD2 microbatch reconciliation. The sort order of affected rows is especially central to how the algorithm works, all other logic follows fairly naturally if we agree on sort order. I put thought into best-effort gracefully handling edge cases like duplicate sequence events (including full row duplicates).

Bundle the (recordStartAt, startAt, endAt) triple that RowClassifier
predicates classify on into a Scd2IntervalColumns case class, and move
the effectiveRecordStartAt coalesce onto it. Converts all RowClassifier
predicates and their call sites to the bundle.
Comment on lines +543 to +601
/**
* Asserts that every row in `decomposedRowsPerKey` conforms to one of the four canonical
* post-decomposition shapes - tombstone, open upsert, closed upsert, or decomposition
* tail - and is otherwise a structural identity transform.
*
* @param decomposedRowsPerKey
* the output of [[decomposeOutOfOrderRows]]: a dataframe conforming to the canonical
* SCD2 row schema `[user_cols..., [[startAtColName]], [[endAtColName]],
* [[cdcMetadataColName]]]`.
* @return
* a dataframe with the exact same schema and rows as the input. Failing the
* well-formedness check is treated as an internal-invariant violation: at execution
* time, the first ill-formed row encountered aborts the query with a SparkRuntimeException
*/
private[autocdc] def assertWellFormedRowsPostDecomposition(
decomposedRowsPerKey: DataFrame,
batchId: Long
): DataFrame = {
val recordStartAtField =
Scd2BatchProcessor.recordStartAtOf(F.col(AutoCdcReservedNames.cdcMetadataColName))
val startAtCol = F.col(Scd2BatchProcessor.startAtColName)
val endAtCol = F.col(Scd2BatchProcessor.endAtColName)

val isWellFormedRow =
RowClassifier.isDecompositionTail(recordStartAtField, startAtCol, endAtCol) ||
RowClassifier.isTombstone(recordStartAtField, startAtCol, endAtCol) ||
RowClassifier.isUpsertRepresentingRow(recordStartAtField, startAtCol, endAtCol)

def stringOrNullLit(c: Column): Column = F.coalesce(c.cast(StringType), F.lit("null"))
val malformedRowDiagnostic = F.concat(
F.lit(
s"During SCD2 reconciliation of microbatch [id=${batchId}], encountered a " +
"post-decomposition row of unexpected shape:"
),
F.lit(s" ${Scd2BatchProcessor.recordStartAtFieldName}="),
stringOrNullLit(recordStartAtField),
F.lit(s", ${Scd2BatchProcessor.startAtColName}="),
stringOrNullLit(startAtCol),
F.lit(s", ${Scd2BatchProcessor.endAtColName}="),
stringOrNullLit(endAtCol),
F.lit(".")
)

val internalErrorOnMalformed = ExpressionUtils.column(
If(
predicate = ExpressionUtils.expression(isWellFormedRow),
trueValue = Literal(null, BooleanType),
falseValue = RaiseError(
Literal("INTERNAL_ERROR"),
CreateMap(Seq(
Literal("message"),
ExpressionUtils.expression(malformedRowDiagnostic)
)),
BooleanType
)
)
)
decomposedRowsPerKey.filter(internalErrorOnMalformed.isNull)
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I agree on balance that this makes sense. This is a compromise being forced by the untyped nature of dataframes, in a more strongly typed system we would 1000% do what's necessary to make this check happen at compile time.

* Synthetic right boundary created by splitting a closed row, temporarily present during
* microbatch reconciliation but never materializes in the target or aux tables.
*/
private[autocdc] def isDecompositionTail(row: Scd2IntervalColumns): Column =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please try converting this into an unapply() case match. I'm not 100% sure that it will work, but my instinct is that it will be much more maintainable that way; right now there's a quite complex implicit invariant that the methods of this classifier cover all the cases and don't overlap.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thought about this.

The biggest paper cut to consider is these classifier functions all accept lazy Column references, and return a lazy Column expression. They must, since dataframe execution itself is lazy.

That means any unapply style function cannot actually know the classification it will return at invocation time, but rather at dataframe execution time.That also means any callers of such a function will still have to multiplex on the returned expression's result at execution time.

The closest thing we can do as per Claude's suggestion is something like the following:

sealed abstract class RowShape(val name: String)
case object DecompositionTail extends RowShape("decomposition_tail")
case object DeleteRow         extends RowShape("delete")
case object OpenUpsert        extends RowShape("open_upsert")
case object ClosedUpsert      extends RowShape("closed_upsert")
case object Malformed         extends RowShape("malformed")
  
def shapeOf(row: Scd2IntervalColumns): Column = {
    val r = row.recordStartAt
    val s = row.startAt
    val e = row.endAt
    F.when(r.isNull && s.isNull && e.isNotNull, DecompositionTail.name)
      .when(r.isNotNull && s.isNotNull && e.isNotNull && s === r && e === r, DeleteRow.name)
      .when(r.isNotNull && s.isNotNull && e.isNull && s <= r, OpenUpsert.name)
      .when(
        r.isNotNull && s.isNotNull && e.isNotNull && s <= r && r < e && s < e,
        ClosedUpsert.name)
      .otherwise(Malformed.name)
  }

But users still have to reason whether each branch is mutually exclusive or not (today it is), and if isn't then they need to reason about the ordering of the conditional branches.

And as mentioned above, callers still need to multiplex on the string returned during dataframe execution - choosing not to do this.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

OK, this makes sense to me.

@AnishMahto
AnishMahto requested a review from jose-torres June 30, 2026 03:02
* Synthetic right boundary created by splitting a closed row, temporarily present during
* microbatch reconciliation but never materializes in the target or aux tables.
*/
private[autocdc] def isDecompositionTail(row: Scd2IntervalColumns): Column =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

OK, this makes sense to me.

jose-torres pushed a commit that referenced this pull request Jul 6, 2026
… rows

Approved AutoCDC SPIP: https://lists.apache.org/thread/j6sj9wo9odgdpgzlxtvhoy7szs0jplf7

--------

### What changes were proposed in this pull request?
**Preamble:**

The SCD type 2 flow is a foreachBatch streaming query on an input change-data-feed, and is responsible for reconciling the incoming change data onto some target table that follows SCD2 replication semantics.

SCD2 flows also maintain an "auxiliary" table to keep track of early-arriving out-of-order received events state. Each microbatch will need to reconcile against this auxiliary table as well, and update the auxiliary table's state appropriately for future microbatches.

**Decompose affected rows:**

Given the set of affected rows in the current microbatch execution - incoming rows in the microbatch, affected rows from aux table, affected rows from target table - the first step in microbatch reconciliation is decomposing closed historical rows that are being bisected by the microbatch.

A closed historical row is a row in the target table that has a non-null start-at and end-at. It's possible an incoming upsert/delete in the microbatch lands with a sequence in between an existing closed row's start/end at (i.e is a late-arriving event), bisecting it.

Decomposing a closed row means exactly this - bisecting the closed interval into a left and right end point, called the decomposed head and tail of the original closed row respectively. The head represents some past upsert event, the tail represents some past delete or upsert-overtake event.

Once a closed row is decomposed into its end points, it can either coalesce with other endpoints/events from the full set of affected rows to form a new historical row, or it can be demoted back to the aux table as a tombstone or no-op upsert.

**Assert well formed rows post-decomposition**
Given the decomposed affected rows, assert each classifies into one of the known row-types; tombstone, synthetic decomposition tail, or upsert representing row (which can be further classified as a valid open or closed upsert).

**Drop redundant rows post-decomposition:**

Decomposition transforms a single row into two synthetic children rows (decomposition head and tail), and its possible the resulting decomposition tail is made logically redundant by the incoming microbatch.

It's also possible that there are duplicate events by key+sequencing, either both introduced in the same microbatch or across the incoming and a past microbatch. Post-decomposition is a good time to reconcile these duplicates, because by this point we have the maximal possible set of rows that will be merged back into the target/aux tables.

We must drop redundant rows prior to reconciling the new start/end-ats for this decomposed set of microbatch-affected rows, to:
1. Prevent zero width rows (start at == end at) from ever appearing in the target table.
2. Ensure reconciliation of a row's start/end at is fully derived from at most one other row in the decomposed set.

### Why are the changes needed?
Needed for core SCD2 reconciliation logic.

### Does this PR introduce _any_ user-facing change?
No, SCD2 is unreleased.

### How was this patch tested?
Unit tests added to `Scd2BatchProcessorSuite`.

### Was this patch authored or co-authored using generative AI tooling?
Co-authored with Claude Opus 4.7

Closes #56311 from AnishMahto/SPARK-57222-SCD2-decompose-affected-rows.

Authored-by: AnishMahto <anish.mahto99@gmail.com>
Signed-off-by: Jose Torres <jtorres@apache.org>
(cherry picked from commit 4933473)
Signed-off-by: Jose Torres <jtorres@apache.org>
@jose-torres

Copy link
Copy Markdown
Contributor

merged to master and 4.x

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants