From a1cff61583f50dbed68dae1ed35fe334b5910f10 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Fri, 18 Sep 2026 16:42:48 +0000 Subject: [PATCH 01/13] [SPARK-59358][PIPELINES] Coalesce ignored null columns during SCD2 reconciliation --- .../autocdc/Scd2BatchProcessor.scala | 410 +++++++++++++++++- .../autocdc/Scd2ForeachBatchHandler.scala | 9 +- .../autocdc/Scd2BatchProcessorSuite.scala | 38 ++ 3 files changed, 451 insertions(+), 6 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 0a83b03c57fe4..855eaec584a8a 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -967,6 +967,162 @@ case class Scd2BatchProcessor( .drop(Scd2BatchProcessor.isRedundantDeleteEncodingColName) } + /** + * Establishes version maps for upsert-representing rows that do not have one. Existing maps are + * preserved, and delete-encoded rows retain null maps because their user-data authorship is not + * meaningful. + */ + private def initializeMissingVersionMaps( + rowsDf: DataFrame, + eligibleSchema: StructType, + ignoreNullSelection: ColumnSelection, + isUpsertRepresentingRow: Column, + resolver: Resolver): DataFrame = { + val metadataColumnName = AutoCdcReservedNames.cdcMetadataColName + val cdcMetadataCol = F.col(metadataColumnName) + val existingVersionMap = Scd2BatchProcessor.versionMapOf(cdcMetadataCol) + val versionMapToPersist = F.when( + isUpsertRepresentingRow && existingVersionMap.isNull, + Scd2VersionMap.buildVersionMap(eligibleSchema, ignoreNullSelection, resolver) + ).otherwise(existingVersionMap) + + rowsDf.withColumn( + metadataColumnName, + cdcMetadataCol + .withField(Scd2BatchProcessor.versionMapFieldName, versionMapToPersist) + .as(metadataColumnName, rowsDf.schema(metadataColumnName).metadata)) + } + + /** + * Replaces unauthored leaf values with values inherited from the most recent authoring row + * for the same key, and materializes version map entries for leaves that would otherwise lose + * their unauthored signal after inheriting a non-null value. + * + * An existing non-null version map is the row's established authorship record. Its semantic + * contents are frozen across ignore-null selection changes: the active selection determines + * which unauthored leaves are coalesced in this reconciliation, but does not reinterpret the + * authorship of any leaf. Removing a leaf from the selection therefore stops coalescing it, + * while adding it begins consulting the authorship already recorded for that leaf. + * + * A null version map is the absence of an authorship record, rather than an established record + * that says every leaf was authored. When ignore-null is enabled, an upsert-representing row + * with a null map lazily establishes one from its stored values and the active selection before + * coalescing. Once established, it follows the same frozen-authorship rule. Schema evolution may + * later materialize an explicit unauthored entry, but does not change a leaf's authorship. + * + * Must run after decomposition cleanup, whose canonical row shapes expose the persisted + * interval gaps and delete boundaries that terminate inheritance, and before start/end + * reconciliation so inherited tracked-history values determine the final SCD2 runs. + */ + private[autocdc] def coalesceIgnoredNulls( + decomposedDf: DataFrame, + ignoreNullSelection: ColumnSelection): DataFrame = { + val resolver = decomposedDf.sparkSession.sessionState.conf.resolver + val userDataColumnSchema = AutoCdcSchemaUtils.excludeColumns( + schema = decomposedDf.schema, + columnNamesToExclude = changeArgs.keys.map(_.name) ++ + Scd2BatchProcessor.reservedFrameworkColNames, + resolver = resolver + ) + val activeIgnoreNullLeafPaths = + Scd2VersionMap.resolveIgnoreNullLeafPaths(userDataColumnSchema, ignoreNullSelection, resolver) + + val cdcMetadataCol = F.col(AutoCdcReservedNames.cdcMetadataColName) + val versionMapCol = Scd2BatchProcessor.versionMapOf(cdcMetadataCol) + + val currentRow = Scd2IntervalColumns( + recordStartAt = Scd2BatchProcessor.recordStartAtOf(cdcMetadataCol), + startAt = F.col(Scd2BatchProcessor.startAtColName), + endAt = F.col(Scd2BatchProcessor.endAtColName) + ) + val nextRow = currentRow.leadBy(1, orderChronologicallyPerKeyWindow) + + val isUpsertRepresentingRow = RowClassifier.isUpsertRepresentingRow(currentRow) + val withInitializedVersionMaps = initializeMissingVersionMaps( + rowsDf = decomposedDf, + eligibleSchema = userDataColumnSchema, + ignoreNullSelection = ignoreNullSelection, + isUpsertRepresentingRow = isUpsertRepresentingRow, + resolver = resolver + ) + + // Strictly preceding rows only, as a row may never inherit from itself. + val precedingRowsInKeyWindow = + orderChronologicallyPerKeyWindow.rowsBetween(Window.unboundedPreceding, -1) + + // Project row-level context used to compute the leaf-level inheritance expressions below. + val (withRowInheritanceContextDf, rowInheritanceContext) = + RowInheritanceContext.projectOn( + df = withInitializedVersionMaps, + rowEndsInheritanceChain = F.coalesce( + // Any row that fully closes or represents a delete event (closing a preceding upsert + // event) ends any running inheritance chain. + RowClassifier.isTombstone(currentRow) || + RowClassifier.isDecompositionTail(currentRow) || + RowClassifier.rowClosesStrictlyBeforeNextRowIsVisible(currentRow.endAt, nextRow), + F.lit(false) + ), + isUpsertRepresentingRow = isUpsertRepresentingRow, + isFirstRowInKeyWindow = + F.row_number().over(orderChronologicallyPerKeyWindow) === 1, + ) + + // Only leaves in the active ignore-null selection participate in reconciliation. The + // initialized version map determines whether each row authored its value; the selection + // determines whether that recorded authorship is relevant to this reconciliation. + val ignoreNullLeafInheritanceContexts = + activeIgnoreNullLeafPaths.zipWithIndex.map { case (path, index) => + LeafInheritanceContext( + path = path, + index = index, + versionMap = versionMapCol, + rowInheritanceContext = rowInheritanceContext, + precedingRowsInKeyWindow = precedingRowsInKeyWindow + ) + } + val ignoreNullContextsByTopLevelColumn = + ignoreNullLeafInheritanceContexts.groupBy(_.path.head) + + // Evaluate all values available to inherit in one window pass. + val withValuesToInheritDf = withRowInheritanceContextDf.withColumns( + ignoreNullLeafInheritanceContexts.map { context => + context.valueToInheritIfAnyColName -> context.valueToInheritIfAny + }.toMap) + + // Compute updated version maps after materializing new entries due to schema evolution. + val schemaEvolutionUpdatedVersionMap = + Scd2BatchProcessor.updateVersionMapWithSchemaEvolution( + versionMapCol, + ignoreNullLeafInheritanceContexts) + + // Build one expression per original column, replacing the CDC metadata and selected user-data + // columns while preserving the original column order. Expressions are resolved against + // columns in [[withValuesToInheritDf]]. + val outputColumns = decomposedDf.columns.map { + case colName if colName == AutoCdcReservedNames.cdcMetadataColName => + // Update the version map after materializing schema-evolution-attributed columns. + cdcMetadataCol + .withField( + Scd2BatchProcessor.versionMapFieldName, + schemaEvolutionUpdatedVersionMap + ) + .as(colName, decomposedDf.schema(colName).metadata) + case colName + if ignoreNullContextsByTopLevelColumn.contains(colName) => + // Rebuild only user-data columns containing an active ignore-null leaf. + val field = userDataColumnSchema(colName) + Scd2BatchProcessor.constructCoalescedIgnoreNullColumn( + Seq(colName), field.dataType, ignoreNullContextsByTopLevelColumn(colName) + ).as(colName, field.metadata) + case colName => + // Pass through keys, framework columns, and user data outside the active selection. + F.col(QuotingUtils.quoteIdentifier(colName)) + } + + // Evaluate the replacements and drop the temporary inheritance columns. + withValuesToInheritDf.select(outputColumns.toImmutableArraySeq: _*) + } + /** * Convert surviving decomposition tails into tombstones. * @@ -1529,6 +1685,62 @@ object Scd2BatchProcessor { .fieldNames .toImmutableArraySeq + /** Materializes explicit unauthored entries required by retroactive schema evolution. */ + private def updateVersionMapWithSchemaEvolution( + versionMap: Column, + leafInheritanceContexts: Seq[LeafInheritanceContext]): Column = { + if (leafInheritanceContexts.isEmpty) { + return versionMap + } + + val conditionalEntries = leafInheritanceContexts.map { context => + F.when( + context.needsSchemaEvolutionEntry, + Scd2VersionMap.buildVersionMapEntry(context.path, authored = false)) + } + val nonNullEntries = F.filter( + F.array(conditionalEntries: _*), (entry: Column) => entry.isNotNull) + val newEntries = F.map_from_entries(nonNullEntries) + F.when( + versionMap.isNotNull, + F.map_concat(versionMap, newEntries)) + } + + /** + * Rebuilds the column at `path` with coalesced leaf values. For struct types the function + * recurses into each field, reassembling the struct from its coalesced children. For leaf + * types it substitutes the inherited value when the leaf's `inherits` predicate holds. + * + * The supplied contexts must be nonempty and cover `path`. Struct fields without a context are + * passed through unchanged. + */ + private def constructCoalescedIgnoreNullColumn( + path: Seq[String], + dataType: DataType, + contextsBeneath: Seq[LeafInheritanceContext]): Column = { + dataType match { + case struct: StructType => + val contextsByChildName = contextsBeneath.groupBy(_.path(path.length)) + val rebuilt = F.struct( + struct.fields.toImmutableArraySeq.map { field => + val childPath = path :+ field.name + contextsByChildName + .get(field.name) + .map(constructCoalescedIgnoreNullColumn(childPath, field.dataType, _)) + .getOrElse(F.col(QuotingUtils.quoteNameParts(childPath))) + .as(field.name) + }: _* + ) + val anyInherits = contextsBeneath.map(_.inherits).reduce(_ || _) + F.when(anyInherits, rebuilt) + .otherwise(F.col(QuotingUtils.quoteNameParts(path))) + case _ => + val context = contextsBeneath.head + F.when(context.inherits, context.valueToInherit) + .otherwise(F.col(QuotingUtils.quoteNameParts(path))) + } + } + /** * Name of temporary column projected onto microbatch to compute the min sequencing value per * key within the microbatch. @@ -1669,11 +1881,9 @@ object Scd2BatchProcessor { // decomposition tails, which are temporarily and synthetically constructed during // reconciliation, have a null record start at. StructField(recordStartAtFieldName, sequencingType, nullable = true), - // The version map representing null-authorship for the row. For persisted rows: - // If the version map is null, that row was ingested with ignore-null off, and all columns - // are considered explicitly authored (null or not). If the version map is non-null, the - // row was ingested with ignore-null on, and contents of the map comply with the contract - // defined in [[Scd2VersionMap]]. + // The version map representing null-authorship for the row. A null map on an + // upsert-representing row means no per-column authorship record has been established yet; + // [[Scd2BatchProcessor.coalesceIgnoredNulls]] may establish one when ignore-null is active. // // Tombstones and decomposition tails also always hold null version maps because column // authorship is not applicable - they are delete markers. @@ -1735,6 +1945,13 @@ private[autocdc] case class Scd2IntervalColumns( F.lead(recordStartAt, offset).over(window), F.lead(startAt, offset).over(window), F.lead(endAt, offset).over(window)) + + /** + * The earliest visible timestamp of this row. For a normal upsert row `startAt` is + * non-null and wins; for a decomposition tail both `recordStartAt` and `startAt` are null + * and `endAt` surfaces via [[effectiveRecordStartAt]]. + */ + def effectiveStartAt: Column = F.coalesce(startAt, effectiveRecordStartAt) } object RowClassifier { @@ -1807,6 +2024,22 @@ object RowClassifier { ): Column = endAt.isNotNull && endAt < nextEffectiveRecordStartAt + /** + * Whether a row closes (`endAt`) strictly before the next row for the same key becomes + * visible, leaving a gap on the visible timeline that only a delete explains. + * + * Distinct from [[rowClosesStrictlyBeforeNextRow]], which compares against the next row's + * event sequence. Before reconciliation a target row standing for a collapsed no-op run + * holds the run's start in `startAt` and its last member's sequence in `recordStartAt`, + * while the members in between sit in the auxiliary table and need not be in the affected + * window. Comparing event sequences would read the interior of such a run as a gap. + */ + private[autocdc] def rowClosesStrictlyBeforeNextRowIsVisible( + endAt: Column, + next: Scd2IntervalColumns + ): Column = + endAt.isNotNull && endAt < next.effectiveStartAt + /** * Whether `row` carries no new information beyond its immediate successor `next` and so * collapses into that successor's run instead of standing as its own visible interval. It is @@ -1828,3 +2061,170 @@ object RowClassifier { !rowClosesStrictlyBeforeNextRow(row.endAt, next.effectiveRecordStartAt) && areTrackedColumnsEqual } + +/** Inheritance properties that apply to a row as a whole rather than to an individual leaf. */ +private[autocdc] case class RowInheritanceContext( + rowEndsInheritanceChain: Column, + isFirstRowInKeyWindow: Column, + isUpsertRepresentingRow: Column) { + + /** Whether the affected window begins with an upsert-representing row. */ + def isFirstUpsertRepresentingRow: Column = + isFirstRowInKeyWindow && isUpsertRepresentingRow +} + +private[autocdc] object RowInheritanceContext { + + /** + * Projects the supplied expressions as temporary DataFrame columns and returns references to + * them. Keeping row-level window expressions out of each leaf's candidate expression prevents + * Catalyst from producing separate copies for every leaf. + */ + def projectOn( + df: DataFrame, + rowEndsInheritanceChain: Column, + isUpsertRepresentingRow: Column, + isFirstRowInKeyWindow: Column): (DataFrame, RowInheritanceContext) = { + val projectedDf = df.withColumns(Map( + rowEndsInheritanceChainColName -> rowEndsInheritanceChain, + isFirstRowInKeyWindowColName -> isFirstRowInKeyWindow, + isUpsertRepresentingRowColName -> isUpsertRepresentingRow + )) + val projectedContext = RowInheritanceContext( + rowEndsInheritanceChain = F.col(rowEndsInheritanceChainColName), + isFirstRowInKeyWindow = F.col(isFirstRowInKeyWindowColName), + isUpsertRepresentingRow = F.col(isUpsertRepresentingRowColName)) + (projectedDf, projectedContext) + } + + private val rowEndsInheritanceChainColName = + s"${AutoCdcReservedNames.prefix}row_ends_inheritance_chain" + + private val isFirstRowInKeyWindowColName = + s"${AutoCdcReservedNames.prefix}is_first_row_in_key_window" + + private val isUpsertRepresentingRowColName = + s"${AutoCdcReservedNames.prefix}is_upsert_representing_row" +} + +/** + * Per-leaf expressions that [[Scd2BatchProcessor.coalesceIgnoredNulls]] evaluates for one + * leaf of the eligible user-data schema. + * + * [[valueToInheritIfAny]] is projected into [[valueToInheritIfAnyColName]] so that + * [[inherits]], [[valueToInherit]], and [[needsSchemaEvolutionEntry]] stay as plain expressions + * over an ordinary column. + * + * @param path the leaf's name parts within the row. + * @param index ordinal used to derive the temporary projected column name. + * @param valueToInheritIfAny latest value contributed by a strictly preceding row, wrapped to + * distinguish no contribution from an explicit null. + * @param inherits whether this row should take [[valueToInherit]] over its stored value. + * @param needsSchemaEvolutionEntry whether this row must gain an explicit unauthored + * entry for the leaf to keep reading as unauthored + * after it inherits. + */ +private[autocdc] case class LeafInheritanceContext( + path: Seq[String], + index: Int, + valueToInheritIfAny: Column, + inherits: Column, + needsSchemaEvolutionEntry: Column) { + + val valueToInheritIfAnyColName: String = + LeafInheritanceContext.valueToInheritIfAnyColumnName(index) + + /** + * Value extracted from the projected wrapper. It is null when the preceding contribution reset + * inheritance or no preceding contribution exists; [[inherits]] distinguishes those cases. + */ + def valueToInherit: Column = + F.col(valueToInheritIfAnyColName) + .getField(LeafInheritanceContext.inheritanceValueFieldName) +} + +private[autocdc] object LeafInheritanceContext { + + /** Builds the per-leaf expressions used to coalesce unauthored nulls from preceding rows. */ + def apply( + path: Seq[String], + index: Int, + versionMap: Column, + rowInheritanceContext: RowInheritanceContext, + precedingRowsInKeyWindow: WindowSpec): LeafInheritanceContext = { + val leafValueInRow = F.col(QuotingUtils.quoteNameParts(path)) + val isLeafAuthoredByRow = + Scd2VersionMap.isAuthored(versionMap, leafValueInRow, path) + + // Each row sees the candidate produced by the latest preceding row that updated the + // inheritance state. The current row can then emit one of three candidate updates for + // next rows: no update, a null update, or a non-null update. A nullable struct + // distinguishes no update from an update whose value is null. + // + // A chain end emits null so later rows cannot inherit across a delete boundary. An authored + // leaf emits its value. If the first row in the affected window is an upsert, it emits its + // stored value even when unauthored because it is the carry-in for omitted earlier history. + val rowResetsInheritanceChain = rowInheritanceContext.rowEndsInheritanceChain + val rowContributesAuthoredValue = + rowInheritanceContext.isUpsertRepresentingRow && isLeafAuthoredByRow + + // When the first row is an upsert, it supplies the baseline inheritance value for every leaf. + // Even if it did not author a leaf itself, its stored value may have been coalesced from rows + // preceding the affected window. + val rowEstablishesCarryIn = rowInheritanceContext.isFirstUpsertRepresentingRow + val rowUpdatesInheritanceChain = + rowResetsInheritanceChain || rowContributesAuthoredValue || rowEstablishesCarryIn + + // A row that does not update the inheritance chain produces a null outer column (hence + // "candidate"). Otherwise, the struct wraps either a null reset or a non-null value; its own + // nullability therefore encodes all three states. + val rowContributedInheritanceCandidate = F.when( + rowUpdatesInheritanceChain, + { + // If the row resets the inheritance chain, then it intentionally updates the inheritance + // chain with a null value. Otherwise it contributes its current leaf value as the new + // inheritance chain value. + val newInheritanceChainValue = F.when(!rowResetsInheritanceChain, leafValueInRow) + F.struct(newInheritanceChainValue.as(inheritanceValueFieldName)) + } + ) + + val valueToInheritIfAnyColName = valueToInheritIfAnyColumnName(index) + val valueToInheritIfAnyColRef = F.col(valueToInheritIfAnyColName) + val precedingRowsContributeInheritableValue = valueToInheritIfAnyColRef.isNotNull + + val valueToInherit = + valueToInheritIfAnyColRef.getField(inheritanceValueFieldName) + // Skip null over null: rebuilding a nested leaf could otherwise turn a null parent struct + // into a non-null struct containing nulls. + val wouldInheritNullOverNull = leafValueInRow.isNull && valueToInherit.isNull + + val rowInheritsLeaf = + rowInheritanceContext.isUpsertRepresentingRow && + !isLeafAuthoredByRow && + precedingRowsContributeInheritableValue && + !wouldInheritNullOverNull + + LeafInheritanceContext( + path = path, + index = index, + valueToInheritIfAny = F.last( + // The value to inherit for any row in the window is the last non-null inheritance + // candidate proposed by a preceding row. Recall if a row contributes a null candidate + // (different from struct{null}) then it declares it has no value to contribute to the + // inheritance chain. + rowContributedInheritanceCandidate, + ignoreNulls = true + ).over(precedingRowsInKeyWindow), + inherits = rowInheritsLeaf, + needsSchemaEvolutionEntry = Scd2VersionMap.needsSchemaEvolutionEntry( + versionMap, leafValueInRow, path, valueToInherit + ) + ) + } + + private val inheritanceValueFieldName: String = "value" + + private def valueToInheritIfAnyColumnName(index: Int): String = + s"${AutoCdcReservedNames.prefix}value_to_inherit_if_any_$index" +} diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala index 124dfba9757ab..5195505e4de84 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala @@ -116,7 +116,14 @@ case class Scd2ForeachBatchHandler( .transform(d => batchProcessor.assertWellFormedRowsPostDecomposition(d, batchId)) .transform(batchProcessor.dropRedundantRowsPostDecomposition) - val reconciledAndRoutedDf = decomposedDf + val withCoalescedIgnoredNullsDf = batchProcessor.changeArgs.ignoreNullSelection match { + case Some(ignoreNullSelection) => + batchProcessor.coalesceIgnoredNulls(decomposedDf, ignoreNullSelection) + case None => + decomposedDf + } + + val reconciledAndRoutedDf = withCoalescedIgnoredNullsDf .transform(batchProcessor.reconcileStartAndEndAt) .transform(batchProcessor.dropLeftoverDeletesPostReconciliation) .transform(batchProcessor.promoteDecompositionTailsToTombstones) diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala index e39de521535df..b2e27577e4fcb 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala @@ -2335,6 +2335,44 @@ class Scd2BatchProcessorSuite extends QueryTest with SharedSparkSession { ) } + // =============== coalesceIgnoredNulls tests =============== + + test("coalesceIgnoredNulls resets stale inherited values after a leading tombstone") { + val ignoreNullSelection = + ColumnSelection.IncludeColumns(Seq(UnqualifiedColumnName("value"))) + val processor = Scd2BatchProcessor( + changeArgs = ChangeArgs( + keys = Seq(UnqualifiedColumnName("id")), + sequencing = F.col("seq"), + storedAsScdType = ScdType.Type2, + ignoreNullSelection = Some(ignoreNullSelection) + ), + resolvedSequencingType = LongType + ) + val userSchema = new StructType().add("id", IntegerType).add("value", StringType) + val valueVersionMapKey = QuotingUtils.quoteNameParts(Seq("value")) + val unauthoredValueVersionMap = Map(valueVersionMapKey -> false) + + // The tombstone is the first row in the affected suffix. The following upserts still carry a + // value inherited before that delete arrived, but their version maps record that they did not + // author it. The tombstone must reset both rows rather than letting the first upsert establish + // the stale value as the window's carry-in. + val df = targetTableOf(userSchema)( + Row(1, null, 20L, 20L, Row(20L, null)), + Row(1, "stale", 10L, null, Row(30L, unauthoredValueVersionMap)), + Row(1, "stale", 10L, null, Row(40L, unauthoredValueVersionMap)) + ) + + checkAnswer( + df = processor.coalesceIgnoredNulls(df, ignoreNullSelection), + expectedAnswer = Seq( + Row(1, null, 20L, 20L, Row(20L, null)), + Row(1, null, 10L, null, Row(30L, unauthoredValueVersionMap)), + Row(1, null, 10L, null, Row(40L, unauthoredValueVersionMap)) + ) + ) + } + // =============== reconcileStartAndEndAt tests =============== test("reconcileStartAndEndAt: a fresh-key run head propagates its startAt to its " + From a44cb86b8dc9fb70096fd169e8e9011f58522b8e Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 21 Sep 2026 16:26:08 +0000 Subject: [PATCH 02/13] add tests --- .../autocdc/Scd2BatchProcessorSuite.scala | 38 -- .../Scd2CoalesceIgnoredNullsSuite.scala | 526 ++++++++++++++++++ .../Scd2ForeachBatchHandlerSuite.scala | 102 +++- 3 files changed, 621 insertions(+), 45 deletions(-) create mode 100644 sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala index b2e27577e4fcb..e39de521535df 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessorSuite.scala @@ -2335,44 +2335,6 @@ class Scd2BatchProcessorSuite extends QueryTest with SharedSparkSession { ) } - // =============== coalesceIgnoredNulls tests =============== - - test("coalesceIgnoredNulls resets stale inherited values after a leading tombstone") { - val ignoreNullSelection = - ColumnSelection.IncludeColumns(Seq(UnqualifiedColumnName("value"))) - val processor = Scd2BatchProcessor( - changeArgs = ChangeArgs( - keys = Seq(UnqualifiedColumnName("id")), - sequencing = F.col("seq"), - storedAsScdType = ScdType.Type2, - ignoreNullSelection = Some(ignoreNullSelection) - ), - resolvedSequencingType = LongType - ) - val userSchema = new StructType().add("id", IntegerType).add("value", StringType) - val valueVersionMapKey = QuotingUtils.quoteNameParts(Seq("value")) - val unauthoredValueVersionMap = Map(valueVersionMapKey -> false) - - // The tombstone is the first row in the affected suffix. The following upserts still carry a - // value inherited before that delete arrived, but their version maps record that they did not - // author it. The tombstone must reset both rows rather than letting the first upsert establish - // the stale value as the window's carry-in. - val df = targetTableOf(userSchema)( - Row(1, null, 20L, 20L, Row(20L, null)), - Row(1, "stale", 10L, null, Row(30L, unauthoredValueVersionMap)), - Row(1, "stale", 10L, null, Row(40L, unauthoredValueVersionMap)) - ) - - checkAnswer( - df = processor.coalesceIgnoredNulls(df, ignoreNullSelection), - expectedAnswer = Seq( - Row(1, null, 20L, 20L, Row(20L, null)), - Row(1, null, 10L, null, Row(30L, unauthoredValueVersionMap)), - Row(1, null, 10L, null, Row(40L, unauthoredValueVersionMap)) - ) - ) - } - // =============== reconcileStartAndEndAt tests =============== test("reconcileStartAndEndAt: a fresh-key run head propagates its startAt to its " + diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala new file mode 100644 index 0000000000000..658455a16c2b0 --- /dev/null +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -0,0 +1,526 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.spark.sql.pipelines.autocdc + +import org.apache.spark.sql.{functions => F, AnalysisException, QueryTest, Row} +import org.apache.spark.sql.catalyst.util.QuotingUtils +import org.apache.spark.sql.classic.DataFrame +import org.apache.spark.sql.internal.SQLConf +import org.apache.spark.sql.test.SharedSparkSession +import org.apache.spark.sql.types._ + +class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { + + private def includeColumns(columnNames: String*): ColumnSelection = + ColumnSelection.IncludeColumns(columnNames.map(UnqualifiedColumnName(_))) + + private def excludeColumns(columnNames: String*): ColumnSelection = + ColumnSelection.ExcludeColumns(columnNames.map(UnqualifiedColumnName(_))) + + private def versionMap(entries: (Seq[String], Boolean)*): Map[String, Boolean] = + entries.map { case (path, authored) => + QuotingUtils.quoteNameParts(path) -> authored + }.toMap + + private def cdcMetadata( + recordStartAt: java.lang.Long, + versionMap: Map[String, Boolean]): Row = + Row(recordStartAt, versionMap) + + private def targetTableOf( + userSchema: StructType, + cdcColumnMetadata: Metadata = Metadata.empty)(rows: Row*): DataFrame = { + val schema = StructType(userSchema.fields.toSeq ++ Seq( + StructField(Scd2BatchProcessor.startAtColName, LongType, nullable = true), + StructField(Scd2BatchProcessor.endAtColName, LongType, nullable = true), + StructField( + AutoCdcReservedNames.cdcMetadataColName, + Scd2BatchProcessor.cdcMetadataColSchema(LongType), + nullable = false, + cdcColumnMetadata) + )) + spark.createDataFrame(spark.sparkContext.parallelize(rows), schema) + } + + private def processor(ignoreNullSelection: ColumnSelection): Scd2BatchProcessor = + Scd2BatchProcessor( + changeArgs = ChangeArgs( + keys = Seq(UnqualifiedColumnName("id")), + sequencing = F.col("seq"), + storedAsScdType = ScdType.Type2, + ignoreNullSelection = Some(ignoreNullSelection) + ), + resolvedSequencingType = LongType + ) + + private def coalesce( + df: DataFrame, + ignoreNullSelection: ColumnSelection): DataFrame = + processor(ignoreNullSelection).coalesceIgnoredNulls(df, ignoreNullSelection) + + test("rows without version maps establish authorship using the current selection") { + val selection = includeColumns("selected") + val schema = new StructType() + .add("id", IntegerType) + .add("selected", StringType) + .add("unselected", StringType) + val existingMap = versionMap(Seq("unselected") -> true) + val initializedMap = versionMap( + Seq("selected") -> false, + Seq("unselected") -> true) + + val input = targetTableOf(schema)( + Row(1, "source", null, 10L, null, cdcMetadata(10L, existingMap)), + Row(1, null, null, 20L, null, cdcMetadata(20L, null)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "source", null, 10L, null, cdcMetadata(10L, existingMap)), + Row(1, "source", null, 20L, null, cdcMetadata(20L, initializedMap)) + ) + ) + } + + test("current selection gates reconciliation while preserving existing map authorship") { + // Existing maps remain the source of truth for ingestion-time authorship. The current + // selection only chooses which leaves to reconcile: adding a leaf consults its recorded + // authorship, while removing a leaf stops coalescing it without rewriting its map entry. + val selection = includeColumns("selectedNow") + val schema = new StructType() + .add("id", IntegerType) + .add("selectedNow", StringType) + .add("removedNow", StringType) + val emptyMap = versionMap() + val authoredSelectedMap = versionMap( + Seq("selectedNow") -> true, + Seq("removedNow") -> false) + val unauthoredSelectedMap = versionMap( + Seq("selectedNow") -> false, + Seq("removedNow") -> false) + + val input = targetTableOf(schema)( + Row(1, "selected-1", "removed-1", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, null, 20L, null, cdcMetadata(20L, authoredSelectedMap)), + Row(2, "selected-2", "removed-2", 10L, null, cdcMetadata(10L, emptyMap)), + Row(2, null, null, 20L, null, cdcMetadata(20L, unauthoredSelectedMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "selected-1", "removed-1", 10L, null, cdcMetadata(10L, emptyMap)), + // The newly selected leaf retains its previously recorded authored null. + Row(1, null, null, 20L, null, cdcMetadata(20L, authoredSelectedMap)), + Row(2, "selected-2", "removed-2", 10L, null, cdcMetadata(10L, emptyMap)), + // The selected leaf inherits; the removed leaf is no longer reconciled. + Row(2, "selected-2", null, 20L, null, cdcMetadata(20L, unauthoredSelectedMap)) + ) + ) + } + + test("schema-evolved leaves gain entries only when they inherit non-null values") { + val selection = includeColumns("evolved") + val schema = new StructType() + .add("id", IntegerType) + .add("evolved", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("evolved") -> false) + + val input = targetTableOf(schema)( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + // No entry models a leaf added after this row's non-null map was established. + Row(1, null, 20L, null, cdcMetadata(20L, emptyMap)), + // Without a preceding value, the absent entry remains the sparse unauthored signal. + Row(2, null, 20L, null, cdcMetadata(20L, emptyMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, "source", 20L, null, cdcMetadata(20L, unauthoredMap)), + Row(2, null, 20L, null, cdcMetadata(20L, emptyMap)) + ) + ) + } + + test("each post-decomposition boundary resets inheritance") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + + val input = targetTableOf(schema)( + // Key 1: a leading tombstone resets values inherited before the affected suffix. + Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, "stale", 10L, null, cdcMetadata(30L, unauthoredMap)), + // Key 2: a decomposition tail represents the same kind of delete boundary. + Row(2, null, null, 20L, cdcMetadata(null, null)), + Row(2, "stale", 10L, null, cdcMetadata(30L, unauthoredMap)), + // Key 3: a closed interval followed by a visibility gap also ends inheritance. + Row(3, "before-gap", 10L, 20L, cdcMetadata(10L, emptyMap)), + Row(3, "stale", 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, null, 10L, null, cdcMetadata(30L, unauthoredMap)), + Row(2, null, null, 20L, cdcMetadata(null, null)), + Row(2, null, 10L, null, cdcMetadata(30L, unauthoredMap)), + Row(3, "before-gap", 10L, 20L, cdcMetadata(10L, emptyMap)), + Row(3, null, 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + ) + } + + test("an authored value restarts inheritance after a reset") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + + val input = targetTableOf(schema)( + Row(1, "old", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, "new", 30L, null, cdcMetadata(30L, emptyMap)), + Row(1, null, 40L, null, cdcMetadata(40L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "old", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, "new", 30L, null, cdcMetadata(30L, emptyMap)), + Row(1, "new", 40L, null, cdcMetadata(40L, unauthoredMap)) + ) + ) + } + + test("unselected data and framework columns retain their values and metadata") { + val selection = includeColumns("selected") + def commentMetadata(comment: String): Metadata = + new MetadataBuilder().putString("comment", comment).build() + + val schema = StructType(Seq( + StructField("id", IntegerType, nullable = true, commentMetadata("key")), + StructField("selected", StringType, nullable = true, commentMetadata("selected")), + StructField("unselected", StringType, nullable = true, commentMetadata("unselected")) + )) + val cdcColumnMetadata = commentMetadata("cdc") + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("selected") -> false) + val input = targetTableOf(schema, cdcColumnMetadata)( + Row(1, "source", "keep-1", 10L, 20L, cdcMetadata(10L, emptyMap)), + Row(1, null, "keep-2", 20L, 60L, cdcMetadata(20L, unauthoredMap)) + ) + + val result = coalesce(input, selection) + + assert(result.schema == input.schema) + checkAnswer( + result, + Seq( + Row(1, "source", "keep-1", 10L, 20L, cdcMetadata(10L, emptyMap)), + Row(1, "source", "keep-2", 20L, 60L, cdcMetadata(20L, unauthoredMap)) + ) + ) + } + + gridTest("column selection honors the configured resolver")( + Seq( + (false, "VALUE"), + (true, "Value") + ) + ) { case (caseSensitive, selectedColumnName) => + withSQLConf(SQLConf.CASE_SENSITIVE.key -> caseSensitive.toString) { + val selection = includeColumns(selectedColumnName) + val schema = new StructType() + .add("id", IntegerType) + .add("Value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("Value") -> false) + val input = targetTableOf(schema)( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, "source", 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + ) + } + } + + test("case-sensitive selection rejects mismatched column casing") { + withSQLConf(SQLConf.CASE_SENSITIVE.key -> "true") { + val selection = includeColumns("VALUE") + val schema = new StructType() + .add("id", IntegerType) + .add("Value", StringType) + val input = targetTableOf(schema)( + Row(1, "source", 10L, null, cdcMetadata(10L, versionMap())) + ) + + val exception = try { + coalesce(input, selection).collect() + throw new IllegalStateException("Expected a case-sensitive column-resolution failure") + } catch { + case e: AnalysisException => e + } + assert(exception.getCondition == "AUTOCDC_COLUMNS_NOT_FOUND_IN_SCHEMA") + } + } + + test("coalescing is idempotent and a previously coalesced first row supplies carry-in") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + val input = targetTableOf(schema)( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + + val firstPass = coalesce(input, selection) + val firstPassRows = Seq( + Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, "source", 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + checkAnswer(firstPass, firstPassRows) + checkAnswer(coalesce(firstPass, selection), firstPassRows) + + val affectedSuffix = firstPass.filter( + F.col(AutoCdcReservedNames.cdcMetadataColName) + .getField(Scd2BatchProcessor.recordStartAtFieldName) >= 20L) + val incoming = targetTableOf(schema)( + Row(1, null, 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + + checkAnswer( + coalesce(affectedSuffix.unionByName(incoming), selection), + Seq( + Row(1, "source", 20L, null, cdcMetadata(20L, unauthoredMap)), + Row(1, "source", 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + ) + } + + test("incoming authored values replace carry-in without unauthored rows interrupting it") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + val input = targetTableOf(schema)( + // The affected suffix starts with an existing, previously coalesced value. + Row(1, "old", 10L, null, cdcMetadata(10L, unauthoredMap)), + Row(1, "new", 20L, null, cdcMetadata(20L, emptyMap)), + // Neither a previously inherited non-null nor an unauthored null contributes a candidate. + Row(1, "stale", 30L, null, cdcMetadata(30L, unauthoredMap)), + Row(1, null, 40L, null, cdcMetadata(40L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "old", 10L, null, cdcMetadata(10L, unauthoredMap)), + Row(1, "new", 20L, null, cdcMetadata(20L, emptyMap)), + Row(1, "new", 30L, null, cdcMetadata(30L, unauthoredMap)), + Row(1, "new", 40L, null, cdcMetadata(40L, unauthoredMap)) + ) + ) + } + + test("an authored null replaces an older non-null inheritance candidate") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val authoredNullMap = versionMap(Seq("value") -> true) + val unauthoredMap = versionMap(Seq("value") -> false) + val input = targetTableOf(schema)( + Row(1, "old", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, null, cdcMetadata(20L, authoredNullMap)), + Row(1, "stale", 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "old", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, 20L, null, cdcMetadata(20L, authoredNullMap)), + Row(1, null, 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + ) + } + + test("each selected leaf maintains an independent inheritance chain") { + val selection = includeColumns("left", "right") + val schema = new StructType() + .add("id", IntegerType) + .add("left", StringType) + .add("right", StringType) + val emptyMap = versionMap() + val leftUnauthoredMap = versionMap(Seq("left") -> false) + val rightUnauthoredMap = versionMap(Seq("right") -> false) + val input = targetTableOf(schema)( + Row(1, "left-1", "right-1", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, "right-2", 20L, null, cdcMetadata(20L, leftUnauthoredMap)), + Row(1, "left-3", null, 30L, null, cdcMetadata(30L, rightUnauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "left-1", "right-1", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, "left-1", "right-2", 20L, null, cdcMetadata(20L, leftUnauthoredMap)), + Row(1, "left-3", "right-2", 30L, null, cdcMetadata(30L, rightUnauthoredMap)) + ) + ) + } + + test("nested leaves inherit without turning null-over-null parents into non-null structs") { + // "Null over null" means both the stored parent struct and the value available to inherit + // are null. The parent must stay null: rebuilding it would turn that null into a non-null + // struct whose fields are all null, even though no non-null value was inherited. + val selection = includeColumns("profile") + val profileSchema = new StructType() + .add("display.name", StringType) + .add("age", IntegerType) + val schema = new StructType() + .add("id", IntegerType) + .add("profile", profileSchema) + .add("outside", StringType) + val authoredAgeMap = versionMap(Seq("profile", "age") -> true) + val authoredProfileMap = versionMap( + Seq("profile", "display.name") -> true, + Seq("profile", "age") -> true) + val unauthoredProfileMap = versionMap( + Seq("profile", "display.name") -> false, + Seq("profile", "age") -> false) + val input = targetTableOf(schema)( + Row(1, Row("Alice", null), "keep-1", 10L, null, cdcMetadata(10L, authoredAgeMap)), + Row(1, null, "keep-2", 20L, null, cdcMetadata(20L, unauthoredProfileMap)), + Row(2, null, "keep-3", 10L, null, cdcMetadata(10L, authoredProfileMap)), + Row(2, null, "keep-4", 20L, null, cdcMetadata(20L, unauthoredProfileMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, Row("Alice", null), "keep-1", 10L, null, + cdcMetadata(10L, authoredAgeMap)), + Row(1, Row("Alice", null), "keep-2", 20L, null, + cdcMetadata(20L, unauthoredProfileMap)), + Row(2, null, "keep-3", 10L, null, cdcMetadata(10L, authoredProfileMap)), + Row(2, null, "keep-4", 20L, null, cdcMetadata(20L, unauthoredProfileMap)) + ) + ) + } + + test("arrays and maps are inherited as opaque leaves") { + val selection = includeColumns("tags", "properties") + val schema = new StructType() + .add("id", IntegerType) + .add("tags", ArrayType(StringType)) + .add("properties", MapType(StringType, IntegerType)) + val emptyMap = versionMap() + val unauthoredMap = versionMap( + Seq("tags") -> false, + Seq("properties") -> false) + val input = targetTableOf(schema)( + Row(1, Seq("one", "two"), Map("x" -> 1), 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, null, null, 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, Seq("one", "two"), Map("x" -> 1), 10L, null, + cdcMetadata(10L, emptyMap)), + Row(1, Seq("one", "two"), Map("x" -> 1), 20L, null, + cdcMetadata(20L, unauthoredMap)) + ) + ) + } + + test("inheritance follows chronological order independently for each key") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + val input = targetTableOf(schema)( + // Physical input order deliberately differs from chronological order and interleaves keys. + Row(1, null, 20L, null, cdcMetadata(20L, unauthoredMap)), + Row(2, "key-2", 5L, null, cdcMetadata(5L, emptyMap)), + Row(1, "key-1", 10L, null, cdcMetadata(10L, emptyMap)), + Row(2, null, 15L, null, cdcMetadata(15L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "key-1", 10L, null, cdcMetadata(10L, emptyMap)), + Row(1, "key-1", 20L, null, cdcMetadata(20L, unauthoredMap)), + Row(2, "key-2", 5L, null, cdcMetadata(5L, emptyMap)), + Row(2, "key-2", 15L, null, cdcMetadata(15L, unauthoredMap)) + ) + ) + } + + test("a selection resolving to no leaves initializes upserts but leaves deletes unchanged") { + val selection = excludeColumns("value", "other") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + .add("other", IntegerType) + val authoredNullsMap = versionMap( + Seq("value") -> true, + Seq("other") -> true) + val input = targetTableOf(schema)( + Row(1, null, null, 10L, null, cdcMetadata(10L, null)), + Row(2, null, null, 20L, 20L, cdcMetadata(20L, null)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, null, null, 10L, null, cdcMetadata(10L, authoredNullsMap)), + Row(2, null, null, 20L, 20L, cdcMetadata(20L, null)) + ) + ) + } +} diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandlerSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandlerSuite.scala index 008313b1ee175..616d5b4f8424e 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandlerSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandlerSuite.scala @@ -20,6 +20,7 @@ package org.apache.spark.sql.pipelines.autocdc import org.scalatest.BeforeAndAfter import org.apache.spark.sql.{functions => F, AnalysisException, QueryTest, Row} +import org.apache.spark.sql.catalyst.util.QuotingUtils import org.apache.spark.sql.classic.{DataFrame, Dataset} import org.apache.spark.sql.internal.SQLConf import org.apache.spark.sql.test.SharedSparkSession @@ -86,6 +87,15 @@ class Scd2ForeachBatchHandlerSuite resolvedSequencingType = LongType ) + private val ignoreNullProcessor = processor.copy( + changeArgs = processor.changeArgs.copy( + ignoreNullSelection = Some(ColumnSelection.IncludeColumns( + Seq(UnqualifiedColumnName("value"))))) + ) + private val emptyVersionMap = Map.empty[String, Boolean] + private val unauthoredValueVersionMap = + Map(QuotingUtils.quoteNameParts(Seq("value")) -> false) + private def createAuxTable(seedRows: Row*): Unit = createTable(defaultAuxIdent, defaultAuxTableIdentifier, auxSchema, seedRows: _*) @@ -120,8 +130,9 @@ class Scd2ForeachBatchHandlerSuite value: String, startAt: java.lang.Long, endAt: java.lang.Long, - recordStartAt: Long): Row = - Row(id, value, startAt, endAt, meta(recordStartAt)) + recordStartAt: Long, + versionMap: Any = null): Row = + Row(id, value, startAt, endAt, meta(recordStartAt, versionMap)) /** A canonical aux row `(id, value, startAt, endAt, meta(recordStartAt), deletedByBatchId)`. */ private def auxRow( @@ -130,13 +141,18 @@ class Scd2ForeachBatchHandlerSuite startAt: java.lang.Long, endAt: java.lang.Long, recordStartAt: Long, - deletedByBatchId: java.lang.Long): Row = - Row(id, value, startAt, endAt, meta(recordStartAt), deletedByBatchId) + deletedByBatchId: java.lang.Long, + versionMap: Any = null): Row = + Row(id, value, startAt, endAt, meta(recordStartAt, versionMap), deletedByBatchId) /** Run a microbatch of source rows through the default handler. */ private def runBatch(batchId: Long)(rows: Row*): Unit = exec.execute(microbatchOf(sourceSchema)(rows: _*), batchId) + /** Run a microbatch with ignore-null enabled for `value`. */ + private def runIgnoreNullBatch(batchId: Long)(rows: Row*): Unit = + execWith(ignoreNullProcessor).execute(microbatchOf(sourceSchema)(rows: _*), batchId) + /** * Run `rows` as batch `batchId`, capture both tables, then replay the identical batch under the * same `batchId` and assert both tables are byte-for-byte unchanged. Models a crash/redelivery @@ -159,11 +175,15 @@ class Scd2ForeachBatchHandlerSuite * updated. On recovery Structured Streaming reruns the same `batchId`, which * [[Scd2BatchProcessor.deletedByBatchIdColName]] is designed to make idempotent. */ - private def runBatchAuxMergeOnly(batchId: Long)(rows: Row*): Unit = { + private def runBatchAuxMergeOnly( + batchId: Long, + p: Scd2BatchProcessor = processor)(rows: Row*): Unit = { // Reuse the handler's own reconciliation chain so this helper cannot drift from execute(), // then run only the aux merge (skipping the target merge) to model the mid-batch crash. - val reconciled = exec.reconcileMicrobatch(microbatchOf(sourceSchema)(rows: _*), batchId) - processor.mergeRowsIntoAuxiliaryTable( + val reconciled = execWith(p).reconcileMicrobatch( + microbatchOf(sourceSchema)(rows: _*), + batchId) + p.mergeRowsIntoAuxiliaryTable( reconciledDfWithAuxRowsTagged = reconciled.reconciledAndRoutedDf, originalAffectedRowsFromAuxiliaryTable = reconciled.affectedRowsFromAuxiliaryTable, auxiliaryTableIdentifier = defaultAuxTableIdentifier, @@ -691,6 +711,74 @@ class Scd2ForeachBatchHandlerSuite checkAnswer(auxTable, auxRow(1, "a", 10L, null, 10L, null)) } + test("recovering after a crash between merges converges with ignore-null coalescing") { + // The null event inherits "a", making both events one run. The first event is therefore + // written to aux before the crash, while the inherited tail and its unauthored map have not + // reached the target. Replaying the same batch must reproduce the clean single-run result. + createAuxTable() + createTargetTable() + + val rows = Seq(upsert(1, "a", 10L), upsert(1, null, 20L)) + runBatchAuxMergeOnly(1L, ignoreNullProcessor)(rows: _*) + + checkAnswer( + auxTable, + auxRow( + 1, "a", 10L, null, 10L, deletedByBatchId = null, versionMap = emptyVersionMap)) + assert(targetTable.collect().isEmpty) + + runIgnoreNullBatch(1L)(rows: _*) + + checkAnswer(targetTable, targetRow(1, "a", 10L, null, 20L, unauthoredValueVersionMap)) + checkAnswer( + auxTable, + auxRow( + 1, "a", 10L, null, 10L, deletedByBatchId = null, versionMap = emptyVersionMap)) + } + + test("retry after an aux-only merge reuses a previously coalesced ignore-null carry-in") { + // Batch 1 establishes a coalesced run. Batch 2 then crashes after demoting its visible tail + // into aux but before advancing the target. The retry sees that row in both tables; it must + // deduplicate the copies, use the coalesced value as carry-in, and advance the run exactly + // once. + createAuxTable() + createTargetTable() + + runIgnoreNullBatch(1L)(upsert(1, "a", 10L), upsert(1, null, 20L)) + checkAnswer(targetTable, targetRow(1, "a", 10L, null, 20L, unauthoredValueVersionMap)) + checkAnswer( + auxTable, + auxRow( + 1, "a", 10L, null, 10L, deletedByBatchId = null, versionMap = emptyVersionMap)) + + runBatchAuxMergeOnly(2L, ignoreNullProcessor)(upsert(1, null, 30L)) + checkAnswer(targetTable, targetRow(1, "a", 10L, null, 20L, unauthoredValueVersionMap)) + checkAnswer( + auxTable, + Seq( + auxRow( + 1, "a", 10L, null, 10L, deletedByBatchId = null, versionMap = emptyVersionMap), + auxRow( + 1, "a", 10L, null, 20L, + deletedByBatchId = null, versionMap = unauthoredValueVersionMap) + ) + ) + + runIgnoreNullBatch(2L)(upsert(1, null, 30L)) + + checkAnswer(targetTable, targetRow(1, "a", 10L, null, 30L, unauthoredValueVersionMap)) + checkAnswer( + auxTable, + Seq( + auxRow( + 1, "a", 10L, null, 10L, deletedByBatchId = null, versionMap = emptyVersionMap), + auxRow( + 1, "a", 10L, null, 20L, + deletedByBatchId = null, versionMap = unauthoredValueVersionMap) + ) + ) + } + test("recovering after a crash that logically deleted a pre-existing aux row converges") { // This is the case the deletedByBatchId re-inclusion clause exists for: unlike the crash tests // above (fresh inserts, deletedByBatchId = null), here the crashed attempt logically DELETES a From 14af9ec2e249d8630b9b4da7bb117f40252b27af Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 21 Sep 2026 16:44:03 +0000 Subject: [PATCH 03/13] explain first upsert row in reconciliation case --- .../spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 855eaec584a8a..69576c3b8901b 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -2171,6 +2171,11 @@ private[autocdc] object LeafInheritanceContext { // When the first row is an upsert, it supplies the baseline inheritance value for every leaf. // Even if it did not author a leaf itself, its stored value may have been coalesced from rows // preceding the affected window. + // + // If this first row is an existing upsert, it cannot itself be re-coalesced in this sweep: + // its predecessor is outside the affected window, so its stored value is the only safe + // carry-in. With an unchanged selection that value is already correct. After a selection + // change, correcting the anchor is deferred until reconciliation includes preceding history. val rowEstablishesCarryIn = rowInheritanceContext.isFirstUpsertRepresentingRow val rowUpdatesInheritanceChain = rowResetsInheritanceChain || rowContributesAuthoredValue || rowEstablishesCarryIn From 058020d1f754e8964907d655323bef7fd0df7e99 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 21 Sep 2026 17:55:56 +0000 Subject: [PATCH 04/13] buff schema evolution test --- .../autocdc/Scd2BatchProcessor.scala | 2 +- .../Scd2CoalesceIgnoredNullsSuite.scala | 22 ++++++++++++++++++- 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 69576c3b8901b..1eb175047aa50 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -1064,7 +1064,7 @@ case class Scd2BatchProcessor( ), isUpsertRepresentingRow = isUpsertRepresentingRow, isFirstRowInKeyWindow = - F.row_number().over(orderChronologicallyPerKeyWindow) === 1, + F.row_number().over(orderChronologicallyPerKeyWindow) === 1 ) // Only leaves in the active ignore-null selection participate in reconciliation. The diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index 658455a16c2b0..d5164ff33c45c 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -151,14 +151,34 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { Row(2, null, 20L, null, cdcMetadata(20L, emptyMap)) ) + val firstPass = coalesce(input, selection) checkAnswer( - coalesce(input, selection), + firstPass, Seq( Row(1, "source", 10L, null, cdcMetadata(10L, emptyMap)), Row(1, "source", 20L, null, cdcMetadata(20L, unauthoredMap)), Row(2, null, 20L, null, cdcMetadata(20L, emptyMap)) ) ) + + // On a later reconciliation, an authored row arrives before the previously inherited row. + // Its explicit false entry keeps the inherited non-null value unauthored, allowing the newer + // preceding value to replace it instead of treating the stored value as authoritative. + val previouslyInheritedRow = firstPass.filter( + F.col("id") === 1 && + F.col(AutoCdcReservedNames.cdcMetadataColName) + .getField(Scd2BatchProcessor.recordStartAtFieldName) === 20L) + val newPrecedingRow = targetTableOf(schema)( + Row(1, "new-source", 15L, null, cdcMetadata(15L, emptyMap)) + ) + + checkAnswer( + coalesce(newPrecedingRow.unionByName(previouslyInheritedRow), selection), + Seq( + Row(1, "new-source", 15L, null, cdcMetadata(15L, emptyMap)), + Row(1, "new-source", 20L, null, cdcMetadata(20L, unauthoredMap)) + ) + ) } test("each post-decomposition boundary resets inheritance") { From 0f5e98f7665d309ae4b60871dfb608b73b2c98cf Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 28 Sep 2026 19:04:33 +0000 Subject: [PATCH 05/13] cleanup scaladoc --- .../spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala | 6 +++--- .../apache/spark/sql/pipelines/autocdc/Scd2VersionMap.scala | 4 ++-- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 1eb175047aa50..568c6ee2cf610 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -968,9 +968,9 @@ case class Scd2BatchProcessor( } /** - * Establishes version maps for upsert-representing rows that do not have one. Existing maps are - * preserved, and delete-encoded rows retain null maps because their user-data authorship is not - * meaningful. + * Establishes version maps from stored values and the current ignore-null selection for + * upsert-representing rows that do not have one. Existing maps are preserved, and delete-encoded + * rows retain null maps because their user-data authorship is not meaningful. */ private def initializeMissingVersionMaps( rowsDf: DataFrame, diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2VersionMap.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2VersionMap.scala index e449ce32f310f..61aed955710e9 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2VersionMap.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2VersionMap.scala @@ -84,8 +84,8 @@ private[pipelines] object Scd2VersionMap { * * @param schema The schema whose leaves the version map covers. Null-authorship is tracked * for every leaf column in this schema, as per the version map contract. - * @param ignoreNullSelection The ignore-null column selection this schema is being ingested - * under. + * @param ignoreNullSelection The ignore-null selection under which to construct the version map + * for the provided schema. * @param resolver Case-sensitivity resolver for column name matching. * @return A [[Column]] of [[mapType]] schema. */ From faebeca4b8fce107143f92305cfe5f8ecd4d9216 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 28 Sep 2026 20:27:37 +0000 Subject: [PATCH 06/13] retain field metadata and type on ignored-null struct reconstruction --- .../autocdc/Scd2BatchProcessor.scala | 5 +- .../Scd2CoalesceIgnoredNullsSuite.scala | 50 +++++++++++++++++++ 2 files changed, 53 insertions(+), 2 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 568c6ee2cf610..1bcf70a4fed69 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -1718,7 +1718,7 @@ object Scd2BatchProcessor { path: Seq[String], dataType: DataType, contextsBeneath: Seq[LeafInheritanceContext]): Column = { - dataType match { + val reconstructed = dataType match { case struct: StructType => val contextsByChildName = contextsBeneath.groupBy(_.path(path.length)) val rebuilt = F.struct( @@ -1728,7 +1728,7 @@ object Scd2BatchProcessor { .get(field.name) .map(constructCoalescedIgnoreNullColumn(childPath, field.dataType, _)) .getOrElse(F.col(QuotingUtils.quoteNameParts(childPath))) - .as(field.name) + .as(field.name, field.metadata) }: _* ) val anyInherits = contextsBeneath.map(_.inherits).reduce(_ || _) @@ -1739,6 +1739,7 @@ object Scd2BatchProcessor { F.when(context.inherits, context.valueToInherit) .otherwise(F.col(QuotingUtils.quoteNameParts(path))) } + reconstructed.cast(dataType) } /** diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index d5164ff33c45c..0b68f999f85a1 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -270,6 +270,56 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { ) } + test("reconstructed nested structs retain schema metadata and nullability") { + def commentMetadata(comment: String): Metadata = + new MetadataBuilder().putString("comment", comment).build() + + val cityMetadata = commentMetadata("city") + val zipMetadata = commentMetadata("zip") + val addressMetadata = commentMetadata("address") + val noteMetadata = commentMetadata("note") + val profileMetadata = commentMetadata("profile") + val addressType = StructType(Seq( + StructField("city", StringType, nullable = false, cityMetadata), + StructField("zip", StringType, nullable = true, zipMetadata) + )) + val profileType = StructType(Seq( + StructField("address", addressType, nullable = false, addressMetadata), + StructField("note", StringType, nullable = true, noteMetadata) + )) + val schema = StructType(Seq( + StructField("id", IntegerType, nullable = false), + StructField("profile", profileType, nullable = false, profileMetadata) + )) + val authoredMap = versionMap( + Seq("profile", "address", "city") -> true, + Seq("profile", "address", "zip") -> true, + Seq("profile", "note") -> true) + val unauthoredZipMap = versionMap( + Seq("profile", "address", "city") -> true, + Seq("profile", "address", "zip") -> false, + Seq("profile", "note") -> true) + val input = targetTableOf(schema)( + Row(1, Row(Row("city-1", "zip-1"), "note-1"), + 10L, null, cdcMetadata(10L, authoredMap)), + Row(1, Row(Row("city-2", null), "note-2"), + 20L, null, cdcMetadata(20L, unauthoredZipMap)) + ) + + val result = coalesce(input, includeColumns("profile")) + + assert(result.schema == input.schema) + checkAnswer( + result, + Seq( + Row(1, Row(Row("city-1", "zip-1"), "note-1"), + 10L, null, cdcMetadata(10L, authoredMap)), + Row(1, Row(Row("city-2", "zip-1"), "note-2"), + 20L, null, cdcMetadata(20L, unauthoredZipMap)) + ) + ) + } + gridTest("column selection honors the configured resolver")( Seq( (false, "VALUE"), From 1a4209b1f2926a9abc77d2f9759dec46ebfaa42b Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Mon, 28 Sep 2026 20:38:59 +0000 Subject: [PATCH 07/13] simplify coalesceIgnoredNulls interface --- .../autocdc/Scd2BatchProcessor.scala | 236 ++++++++++-------- .../autocdc/Scd2ForeachBatchHandler.scala | 8 +- .../Scd2CoalesceIgnoredNullsSuite.scala | 2 +- 3 files changed, 130 insertions(+), 116 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 1bcf70a4fed69..35d3226d2d474 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -23,6 +23,7 @@ import org.apache.spark.sql.Column import org.apache.spark.sql.catalyst.TableIdentifier import org.apache.spark.sql.catalyst.analysis.Resolver import org.apache.spark.sql.catalyst.expressions.{CreateMap, If, Literal, RaiseError} +import org.apache.spark.sql.catalyst.expressions.objects.AssertNotNull import org.apache.spark.sql.catalyst.util.QuotingUtils import org.apache.spark.sql.classic.{DataFrame, ExpressionUtils} import org.apache.spark.sql.expressions.{Window, WindowSpec} @@ -1015,113 +1016,123 @@ case class Scd2BatchProcessor( * reconciliation so inherited tracked-history values determine the final SCD2 runs. */ private[autocdc] def coalesceIgnoredNulls( - decomposedDf: DataFrame, - ignoreNullSelection: ColumnSelection): DataFrame = { - val resolver = decomposedDf.sparkSession.sessionState.conf.resolver - val userDataColumnSchema = AutoCdcSchemaUtils.excludeColumns( - schema = decomposedDf.schema, - columnNamesToExclude = changeArgs.keys.map(_.name) ++ - Scd2BatchProcessor.reservedFrameworkColNames, - resolver = resolver - ) - val activeIgnoreNullLeafPaths = - Scd2VersionMap.resolveIgnoreNullLeafPaths(userDataColumnSchema, ignoreNullSelection, resolver) - - val cdcMetadataCol = F.col(AutoCdcReservedNames.cdcMetadataColName) - val versionMapCol = Scd2BatchProcessor.versionMapOf(cdcMetadataCol) + decomposedDf: DataFrame): DataFrame = + changeArgs.ignoreNullSelection match { + case None => decomposedDf + case Some(ignoreNullSelection) => + val resolver = decomposedDf.sparkSession.sessionState.conf.resolver + val userDataColumnSchema = AutoCdcSchemaUtils.excludeColumns( + schema = decomposedDf.schema, + columnNamesToExclude = changeArgs.keys.map(_.name) ++ + Scd2BatchProcessor.reservedFrameworkColNames, + resolver = resolver + ) + val activeIgnoreNullLeafPaths = + Scd2VersionMap.resolveIgnoreNullLeafPaths( + userDataColumnSchema, ignoreNullSelection, resolver) - val currentRow = Scd2IntervalColumns( - recordStartAt = Scd2BatchProcessor.recordStartAtOf(cdcMetadataCol), - startAt = F.col(Scd2BatchProcessor.startAtColName), - endAt = F.col(Scd2BatchProcessor.endAtColName) - ) - val nextRow = currentRow.leadBy(1, orderChronologicallyPerKeyWindow) - - val isUpsertRepresentingRow = RowClassifier.isUpsertRepresentingRow(currentRow) - val withInitializedVersionMaps = initializeMissingVersionMaps( - rowsDf = decomposedDf, - eligibleSchema = userDataColumnSchema, - ignoreNullSelection = ignoreNullSelection, - isUpsertRepresentingRow = isUpsertRepresentingRow, - resolver = resolver - ) + val cdcMetadataCol = F.col(AutoCdcReservedNames.cdcMetadataColName) + val versionMapCol = Scd2BatchProcessor.versionMapOf(cdcMetadataCol) - // Strictly preceding rows only, as a row may never inherit from itself. - val precedingRowsInKeyWindow = - orderChronologicallyPerKeyWindow.rowsBetween(Window.unboundedPreceding, -1) - - // Project row-level context used to compute the leaf-level inheritance expressions below. - val (withRowInheritanceContextDf, rowInheritanceContext) = - RowInheritanceContext.projectOn( - df = withInitializedVersionMaps, - rowEndsInheritanceChain = F.coalesce( - // Any row that fully closes or represents a delete event (closing a preceding upsert - // event) ends any running inheritance chain. - RowClassifier.isTombstone(currentRow) || - RowClassifier.isDecompositionTail(currentRow) || - RowClassifier.rowClosesStrictlyBeforeNextRowIsVisible(currentRow.endAt, nextRow), - F.lit(false) - ), - isUpsertRepresentingRow = isUpsertRepresentingRow, - isFirstRowInKeyWindow = - F.row_number().over(orderChronologicallyPerKeyWindow) === 1 - ) + val currentRow = Scd2IntervalColumns( + recordStartAt = Scd2BatchProcessor.recordStartAtOf(cdcMetadataCol), + startAt = F.col(Scd2BatchProcessor.startAtColName), + endAt = F.col(Scd2BatchProcessor.endAtColName) + ) + val nextRow = currentRow.leadBy(1, orderChronologicallyPerKeyWindow) - // Only leaves in the active ignore-null selection participate in reconciliation. The - // initialized version map determines whether each row authored its value; the selection - // determines whether that recorded authorship is relevant to this reconciliation. - val ignoreNullLeafInheritanceContexts = - activeIgnoreNullLeafPaths.zipWithIndex.map { case (path, index) => - LeafInheritanceContext( - path = path, - index = index, - versionMap = versionMapCol, - rowInheritanceContext = rowInheritanceContext, - precedingRowsInKeyWindow = precedingRowsInKeyWindow + val isUpsertRepresentingRow = RowClassifier.isUpsertRepresentingRow(currentRow) + val withInitializedVersionMaps = initializeMissingVersionMaps( + rowsDf = decomposedDf, + eligibleSchema = userDataColumnSchema, + ignoreNullSelection = ignoreNullSelection, + isUpsertRepresentingRow = isUpsertRepresentingRow, + resolver = resolver ) - } - val ignoreNullContextsByTopLevelColumn = - ignoreNullLeafInheritanceContexts.groupBy(_.path.head) - - // Evaluate all values available to inherit in one window pass. - val withValuesToInheritDf = withRowInheritanceContextDf.withColumns( - ignoreNullLeafInheritanceContexts.map { context => - context.valueToInheritIfAnyColName -> context.valueToInheritIfAny - }.toMap) - - // Compute updated version maps after materializing new entries due to schema evolution. - val schemaEvolutionUpdatedVersionMap = - Scd2BatchProcessor.updateVersionMapWithSchemaEvolution( - versionMapCol, - ignoreNullLeafInheritanceContexts) - - // Build one expression per original column, replacing the CDC metadata and selected user-data - // columns while preserving the original column order. Expressions are resolved against - // columns in [[withValuesToInheritDf]]. - val outputColumns = decomposedDf.columns.map { - case colName if colName == AutoCdcReservedNames.cdcMetadataColName => - // Update the version map after materializing schema-evolution-attributed columns. - cdcMetadataCol - .withField( - Scd2BatchProcessor.versionMapFieldName, - schemaEvolutionUpdatedVersionMap + + // Strictly preceding rows only, as a row may never inherit from itself. + val precedingRowsInKeyWindow = + orderChronologicallyPerKeyWindow.rowsBetween(Window.unboundedPreceding, -1) + + // Project row-level context used to compute the leaf-level inheritance + // expressions below. + val (withRowInheritanceContextDf, rowInheritanceContext) = + RowInheritanceContext.projectOn( + df = withInitializedVersionMaps, + rowEndsInheritanceChain = F.coalesce( + // Any row that fully closes or represents a delete event (closing a + // preceding upsert event) ends any running inheritance chain. + RowClassifier.isTombstone(currentRow) || + RowClassifier.isDecompositionTail(currentRow) || + RowClassifier.rowClosesStrictlyBeforeNextRowIsVisible( + currentRow.endAt, nextRow), + F.lit(false) + ), + isUpsertRepresentingRow = isUpsertRepresentingRow, + isFirstRowInKeyWindow = + F.row_number().over(orderChronologicallyPerKeyWindow) === 1 ) - .as(colName, decomposedDf.schema(colName).metadata) - case colName - if ignoreNullContextsByTopLevelColumn.contains(colName) => - // Rebuild only user-data columns containing an active ignore-null leaf. - val field = userDataColumnSchema(colName) - Scd2BatchProcessor.constructCoalescedIgnoreNullColumn( - Seq(colName), field.dataType, ignoreNullContextsByTopLevelColumn(colName) - ).as(colName, field.metadata) - case colName => - // Pass through keys, framework columns, and user data outside the active selection. - F.col(QuotingUtils.quoteIdentifier(colName)) - } - // Evaluate the replacements and drop the temporary inheritance columns. - withValuesToInheritDf.select(outputColumns.toImmutableArraySeq: _*) - } + // Only leaves in the active ignore-null selection participate in + // reconciliation. The initialized version map determines whether each row + // authored its value; the selection determines whether that recorded + // authorship is relevant to this reconciliation. + val ignoreNullLeafInheritanceContexts = + activeIgnoreNullLeafPaths.zipWithIndex.map { case (path, index) => + LeafInheritanceContext( + path = path, + index = index, + versionMap = versionMapCol, + rowInheritanceContext = rowInheritanceContext, + precedingRowsInKeyWindow = precedingRowsInKeyWindow + ) + } + val ignoreNullContextsByTopLevelColumn = + ignoreNullLeafInheritanceContexts.groupBy(_.path.head) + + // Evaluate all values available to inherit in one window pass. + val withValuesToInheritDf = withRowInheritanceContextDf.withColumns( + ignoreNullLeafInheritanceContexts.map { context => + context.valueToInheritIfAnyColName -> context.valueToInheritIfAny + }.toMap) + + // Compute updated version maps after materializing new entries due to + // schema evolution. + val schemaEvolutionUpdatedVersionMap = + Scd2BatchProcessor.updateVersionMapWithSchemaEvolution( + versionMapCol, + ignoreNullLeafInheritanceContexts) + + // Build one expression per original column, replacing the CDC metadata and + // selected user-data columns while preserving the original column order. + // Expressions are resolved against columns in [[withValuesToInheritDf]]. + val outputColumns = decomposedDf.columns.map { + case colName if colName == AutoCdcReservedNames.cdcMetadataColName => + // Update the version map after materializing schema-evolution entries. + cdcMetadataCol + .withField( + Scd2BatchProcessor.versionMapFieldName, + schemaEvolutionUpdatedVersionMap + ) + .as(colName, decomposedDf.schema(colName).metadata) + case colName + if ignoreNullContextsByTopLevelColumn.contains(colName) => + // Rebuild only user-data columns containing an active ignore-null leaf. + val field = userDataColumnSchema(colName) + Scd2BatchProcessor.constructCoalescedIgnoreNullColumn( + Seq(colName), + field, + ignoreNullContextsByTopLevelColumn(colName) + ).as(colName, field.metadata) + case colName => + // Pass through keys, framework columns, and user data outside the + // active selection. + F.col(QuotingUtils.quoteIdentifier(colName)) + } + + // Evaluate the replacements and drop the temporary inheritance columns. + withValuesToInheritDf.select(outputColumns.toImmutableArraySeq: _*) + } /** * Convert surviving decomposition tails into tombstones. @@ -1716,19 +1727,19 @@ object Scd2BatchProcessor { */ private def constructCoalescedIgnoreNullColumn( path: Seq[String], - dataType: DataType, + field: StructField, contextsBeneath: Seq[LeafInheritanceContext]): Column = { - val reconstructed = dataType match { + val reconstructed = field.dataType match { case struct: StructType => val contextsByChildName = contextsBeneath.groupBy(_.path(path.length)) val rebuilt = F.struct( - struct.fields.toImmutableArraySeq.map { field => - val childPath = path :+ field.name + struct.fields.toImmutableArraySeq.map { childField => + val childPath = path :+ childField.name contextsByChildName - .get(field.name) - .map(constructCoalescedIgnoreNullColumn(childPath, field.dataType, _)) + .get(childField.name) + .map(constructCoalescedIgnoreNullColumn(childPath, childField, _)) .getOrElse(F.col(QuotingUtils.quoteNameParts(childPath))) - .as(field.name, field.metadata) + .as(childField.name, childField.metadata) }: _* ) val anyInherits = contextsBeneath.map(_.inherits).reduce(_ || _) @@ -1739,7 +1750,14 @@ object Scd2BatchProcessor { F.when(context.inherits, context.valueToInherit) .otherwise(F.col(QuotingUtils.quoteNameParts(path))) } - reconstructed.cast(dataType) + val nullabilityChecked = + if (field.nullable) { + reconstructed + } else { + ExpressionUtils.column( + AssertNotNull(ExpressionUtils.expression(reconstructed), path)) + } + nullabilityChecked.cast(field.dataType) } /** diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala index 5195505e4de84..87ad28b959aa2 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2ForeachBatchHandler.scala @@ -116,12 +116,8 @@ case class Scd2ForeachBatchHandler( .transform(d => batchProcessor.assertWellFormedRowsPostDecomposition(d, batchId)) .transform(batchProcessor.dropRedundantRowsPostDecomposition) - val withCoalescedIgnoredNullsDf = batchProcessor.changeArgs.ignoreNullSelection match { - case Some(ignoreNullSelection) => - batchProcessor.coalesceIgnoredNulls(decomposedDf, ignoreNullSelection) - case None => - decomposedDf - } + val withCoalescedIgnoredNullsDf = + batchProcessor.coalesceIgnoredNulls(decomposedDf) val reconciledAndRoutedDf = withCoalescedIgnoredNullsDf .transform(batchProcessor.reconcileStartAndEndAt) diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index 0b68f999f85a1..f8617b78489fe 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -71,7 +71,7 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { private def coalesce( df: DataFrame, ignoreNullSelection: ColumnSelection): DataFrame = - processor(ignoreNullSelection).coalesceIgnoredNulls(df, ignoreNullSelection) + processor(ignoreNullSelection).coalesceIgnoredNulls(df) test("rows without version maps establish authorship using the current selection") { val selection = includeColumns("selected") From 49b9c594c23ff333a8fa6c4481cc8f048996ab61 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 04:43:14 +0000 Subject: [PATCH 08/13] document eventual consistency behavior/contract --- .../autocdc/Scd2BatchProcessor.scala | 30 +++++- .../Scd2CoalesceIgnoredNullsSuite.scala | 93 +++++++++++++++++++ 2 files changed, 119 insertions(+), 4 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 35d3226d2d474..965f7f7cf3975 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -995,9 +995,10 @@ case class Scd2BatchProcessor( } /** - * Replaces unauthored leaf values with values inherited from the most recent authoring row - * for the same key, and materializes version map entries for leaves that would otherwise lose - * their unauthored signal after inheriting a non-null value. + * Replaces unauthored leaf values with values inherited from the nearest preceding authoring + * row for the same key in `decomposedDf`, and materializes version map entries for + * schema-evolved leaves that would otherwise lose their unauthored signal after inheriting a + * non-null value. * * An existing non-null version map is the row's established authorship record. Its semantic * contents are frozen across ignore-null selection changes: the active selection determines @@ -1014,6 +1015,25 @@ case class Scd2BatchProcessor( * Must run after decomposition cleanup, whose canonical row shapes expose the persisted * interval gaps and delete boundaries that terminate inheritance, and before start/end * reconciliation so inherited tracked-history values determine the final SCD2 runs. + * + * ==Eventual consistency of the carry-in anchor== + * + * A leaf's reconciliation in any pass is a function of two independent inputs: + * + * 1. The active ignore-null selection. Only leaves in the active selection are eligible for + * coalescing; an unauthored leaf outside it is not reconciled in this pass. + * 2. The affected window of rows pulled in for reconciliation. It is computed independently + * of the selection, from anchor-row deduction (the earliest existing row bisected by an + * incoming event, or the earliest event in the microbatch). + * + * The first upsert-representing row in the affected window supplies its stored value as + * carry-in, even when its version map records that value as unauthored, because its + * predecessor is outside the window. If the leaf was temporarily dropped from the selection, + * that stored value may not have been reconciled against the latest authoring row. + * + * Correction for that anchor row is eventually consistent: once the affected window expands to + * include a preceding row that authored the leaf, the stale carry-in is replaced and all + * downstream unauthored rows inherit the correct value. */ private[autocdc] def coalesceIgnoredNulls( decomposedDf: DataFrame): DataFrame = @@ -2194,7 +2214,9 @@ private[autocdc] object LeafInheritanceContext { // If this first row is an existing upsert, it cannot itself be re-coalesced in this sweep: // its predecessor is outside the affected window, so its stored value is the only safe // carry-in. With an unchanged selection that value is already correct. After a selection - // change, correcting the anchor is deferred until reconciliation includes preceding history. + // change, correcting the anchor is deferred until reconciliation includes preceding + // history. See the "Eventual consistency of the carry-in anchor" section in the + // [[coalesceIgnoredNulls]] scaladoc. val rowEstablishesCarryIn = rowInheritanceContext.isFirstUpsertRepresentingRow val rowUpdatesInheritanceChain = rowResetsInheritanceChain || rowContributesAuthoredValue || rowEstablishesCarryIn diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index f8617b78489fe..f6c6f823a61f1 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -593,4 +593,97 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { ) ) } + + test("coalescing is eventually consistent for existing rows when " + + "ignore-null selection changes") { + // Demonstrates the "Eventual consistency of the carry-in anchor" contract of + // [[Scd2BatchProcessor.coalesceIgnoredNulls]]: a leaf temporarily dropped from the ignore-null + // selection is not reconciled, so a later narrow affected window propagates a stale carry-in + // until a wider window pulls in the authoring row. + + // "other" is authored and non-null in every row, so coalescing never changes it. It exists + // only so the selection can move off "value" without becoming an (invalid) empty list. + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + .add("other", StringType) + val emptyMap = versionMap() + val unauthoredMap = versionMap(Seq("value") -> false) + + // -- Batch 1 (selection = {value}) ------------------------------------------------ + // Two rows arrive for key 1: seq 1 authors "old", seq 3 has an unauthored null. + // The affected window spans both rows, so seq 3 inherits "old" from seq 1. + val batch1Input = targetTableOf(schema)( + Row(1, "old", "x", 1L, null, cdcMetadata(1L, emptyMap)), + Row(1, null, "x", 3L, null, cdcMetadata(3L, unauthoredMap)) + ) + checkAnswer( + coalesce(batch1Input, includeColumns("value")), + Seq( + Row(1, "old", "x", 1L, null, cdcMetadata(1L, emptyMap)), + Row(1, "old", "x", 3L, null, cdcMetadata(3L, unauthoredMap)) + ) + ) + + // -- Batch 2 (selection = {other} -- "value" temporarily removed) ----------------- + // An out-of-order seq 2 arrives with value = "new", slotting between seq 1 and + // seq 3. The decomposed target state after batch 2: + // seq 1: value="old", startAt=1, endAt=2, versionMap={} + // seq 2: value="new", startAt=2, endAt=3, versionMap={} + // seq 3: value="old", startAt=3, endAt=null, versionMap={value->false} + val batch2AffectedWindow = targetTableOf(schema)( + Row(1, "old", "x", 1L, 2L, cdcMetadata(1L, emptyMap)), + Row(1, "new", "x", 2L, 3L, cdcMetadata(2L, emptyMap)), + Row(1, "old", "x", 3L, null, cdcMetadata(3L, unauthoredMap)) + ) + // "value" is not in the active selection, so it is not reconciled. Seq 3 retains its + // previously inherited "old" even though seq 2 (which authors "new") now directly + // precedes it. + checkAnswer( + coalesce(batch2AffectedWindow, includeColumns("other")), + Seq( + Row(1, "old", "x", 1L, 2L, cdcMetadata(1L, emptyMap)), + Row(1, "new", "x", 2L, 3L, cdcMetadata(2L, emptyMap)), + Row(1, "old", "x", 3L, null, cdcMetadata(3L, unauthoredMap)) + ) + ) + + // -- Batch 3 (selection = {value} -- "value" re-added) ---------------------------- + // Seq 4 arrives with an unauthored null. The affected window covers only the rows + // touched by this event: seq 3 (the anchor) and seq 4 (the new row). Seq 2, which + // holds the correct "new", is outside the window. + // + // The carry-in anchor (seq 3) has a stale value: it inherited "old" in batch 1 + // and was not reconciled while "value" was outside the selection. Seq 4 inherits + // this stale "old" rather than the correct "new". + val batch3Input = targetTableOf(schema)( + Row(1, "old", "x", 3L, 4L, cdcMetadata(3L, unauthoredMap)), + Row(1, null, "x", 4L, null, cdcMetadata(4L, unauthoredMap)) + ) + checkAnswer( + coalesce(batch3Input, includeColumns("value")), + Seq( + Row(1, "old", "x", 3L, 4L, cdcMetadata(3L, unauthoredMap)), + Row(1, "old", "x", 4L, null, cdcMetadata(4L, unauthoredMap)) + ) + ) + + // -- Correction ------------------------------------------------------------------- + // A future batch's affected window expands to include seq 2. Now seq 2 is the + // first authored row in the window: its "new" propagates forward and the stale + // carry-in is corrected. + val correctionInput = targetTableOf(schema)( + Row(1, "new", "x", 2L, 3L, cdcMetadata(2L, emptyMap)), + Row(1, "old", "x", 3L, 4L, cdcMetadata(3L, unauthoredMap)), + Row(1, "old", "x", 4L, null, cdcMetadata(4L, unauthoredMap)) + ) + checkAnswer( + coalesce(correctionInput, includeColumns("value")), + Seq( + Row(1, "new", "x", 2L, 3L, cdcMetadata(2L, emptyMap)), + Row(1, "new", "x", 3L, 4L, cdcMetadata(3L, unauthoredMap)), + Row(1, "new", "x", 4L, null, cdcMetadata(4L, unauthoredMap)) + ) + ) + } } From 3b7a21d11b98e7d7410ea07e62ecc445a2b08c23 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 04:47:45 +0000 Subject: [PATCH 09/13] cleanup tests --- .../Scd2CoalesceIgnoredNullsSuite.scala | 30 ++++++++++++------- 1 file changed, 19 insertions(+), 11 deletions(-) diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index f6c6f823a61f1..38678502551a7 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -189,12 +189,14 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { val emptyMap = versionMap() val unauthoredMap = versionMap(Seq("value") -> false) + // Every boundary row carries non-null data, so a leak into the inheritance chain would + // surface as that value instead of the expected null. val input = targetTableOf(schema)( // Key 1: a leading tombstone resets values inherited before the affected suffix. - Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, "tombstone", 20L, 20L, cdcMetadata(20L, null)), Row(1, "stale", 10L, null, cdcMetadata(30L, unauthoredMap)), // Key 2: a decomposition tail represents the same kind of delete boundary. - Row(2, null, null, 20L, cdcMetadata(null, null)), + Row(2, "tail", null, 20L, cdcMetadata(null, null)), Row(2, "stale", 10L, null, cdcMetadata(30L, unauthoredMap)), // Key 3: a closed interval followed by a visibility gap also ends inheritance. Row(3, "before-gap", 10L, 20L, cdcMetadata(10L, emptyMap)), @@ -204,9 +206,9 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { checkAnswer( coalesce(input, selection), Seq( - Row(1, null, 20L, 20L, cdcMetadata(20L, null)), + Row(1, "tombstone", 20L, 20L, cdcMetadata(20L, null)), Row(1, null, 10L, null, cdcMetadata(30L, unauthoredMap)), - Row(2, null, null, 20L, cdcMetadata(null, null)), + Row(2, "tail", null, 20L, cdcMetadata(null, null)), Row(2, null, 10L, null, cdcMetadata(30L, unauthoredMap)), Row(3, "before-gap", 10L, 20L, cdcMetadata(10L, emptyMap)), Row(3, null, 30L, null, cdcMetadata(30L, unauthoredMap)) @@ -358,13 +360,19 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { Row(1, "source", 10L, null, cdcMetadata(10L, versionMap())) ) - val exception = try { - coalesce(input, selection).collect() - throw new IllegalStateException("Expected a case-sensitive column-resolution failure") - } catch { - case e: AnalysisException => e - } - assert(exception.getCondition == "AUTOCDC_COLUMNS_NOT_FOUND_IN_SCHEMA") + checkError( + exception = intercept[AnalysisException] { + coalesce(input, selection) + }, + condition = "AUTOCDC_COLUMNS_NOT_FOUND_IN_SCHEMA", + sqlState = "42703", + parameters = Map( + "caseSensitivity" -> CaseSensitivityLabels.CaseSensitive, + "schemaName" -> "ignoreNullSelection", + "missingColumns" -> "VALUE", + "availableColumns" -> "Value" + ) + ) } } From d99f51616384bf855b60fab7762473766a2581fe Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 05:07:04 +0000 Subject: [PATCH 10/13] better document constructCoalescedIgnoreNullColumn --- .../pipelines/autocdc/Scd2BatchProcessor.scala | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 965f7f7cf3975..5b02e71b575d3 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -1751,6 +1751,8 @@ object Scd2BatchProcessor { contextsBeneath: Seq[LeafInheritanceContext]): Column = { val reconstructed = field.dataType match { case struct: StructType => + // If this field is a struct, recursively reconstruct all of its children by applying their + // resolved values after ignore-null coalescing. val contextsByChildName = contextsBeneath.groupBy(_.path(path.length)) val rebuilt = F.struct( struct.fields.toImmutableArraySeq.map { childField => @@ -1762,22 +1764,31 @@ object Scd2BatchProcessor { .as(childField.name, childField.metadata) }: _* ) + // If none of this struct's leaves are inheriting as part of this coalesce pass, let the + // struct pass its value through as-is. This is to avoid incorrectly materializing a null + // struct with a struct with all null leaves. val anyInherits = contextsBeneath.map(_.inherits).reduce(_ || _) F.when(anyInherits, rebuilt) .otherwise(F.col(QuotingUtils.quoteNameParts(path))) case _ => + // If this field is not a struct (and therefore must be a leaf), either directly apply the + // resolved value to inherit if the leaf should be inheriting, otherwise pass its value + // through as-is. val context = contextsBeneath.head F.when(context.inherits, context.valueToInherit) .otherwise(F.col(QuotingUtils.quoteNameParts(path))) } - val nullabilityChecked = + val validatedReconstructed = if (field.nullable) { reconstructed } else { + // If the field was marked as non-nullable but coalescing deduces it will resolve to a + // null, throw an explicit exception. Pushing `AssertNotNull` into the plan for a + // non-nullable field also prevents Spark from preemptively complaining during analysis. ExpressionUtils.column( AssertNotNull(ExpressionUtils.expression(reconstructed), path)) } - nullabilityChecked.cast(field.dataType) + validatedReconstructed.cast(field.dataType) } /** From bdfb6f633b4a37c75277474e662c29bd6268fe2c Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 16:51:14 +0000 Subject: [PATCH 11/13] local PR review --- .../autocdc/Scd2BatchProcessor.scala | 31 ++++++----- .../Scd2CoalesceIgnoredNullsSuite.scala | 54 +++++++++++++------ 2 files changed, 57 insertions(+), 28 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 5b02e71b575d3..05c6ee52353f8 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -996,9 +996,9 @@ case class Scd2BatchProcessor( /** * Replaces unauthored leaf values with values inherited from the nearest preceding authoring - * row for the same key in `decomposedDf`, and materializes version map entries for - * schema-evolved leaves that would otherwise lose their unauthored signal after inheriting a - * non-null value. + * row for the same key in `decomposedDf` (or, absent one, from the carry-in anchor described + * below), and materializes version map entries for schema-evolved leaves that would otherwise + * lose their unauthored signal after inheriting a non-null value. * * An existing non-null version map is the row's established authorship record. Its semantic * contents are frozen across ignore-null selection changes: the active selection determines @@ -1022,18 +1022,19 @@ case class Scd2BatchProcessor( * * 1. The active ignore-null selection. Only leaves in the active selection are eligible for * coalescing; an unauthored leaf outside it is not reconciled in this pass. - * 2. The affected window of rows pulled in for reconciliation. It is computed independently - * of the selection, from anchor-row deduction (the earliest existing row bisected by an - * incoming event, or the earliest event in the microbatch). + * 2. The affected window of rows pulled in for reconciliation, which is computed independently + * of the ignore-null selection. * * The first upsert-representing row in the affected window supplies its stored value as * carry-in, even when its version map records that value as unauthored, because its * predecessor is outside the window. If the leaf was temporarily dropped from the selection, * that stored value may not have been reconciled against the latest authoring row. * - * Correction for that anchor row is eventually consistent: once the affected window expands to - * include a preceding row that authored the leaf, the stale carry-in is replaced and all - * downstream unauthored rows inherit the correct value. + * The stale carry-in is eventually consistent: a later reconciliation replaces it, and all + * downstream unauthored rows inherit the correct value, once that reconciliation's affected + * window includes a preceding row that authored the leaf. That is the only condition for + * correction, and nothing forces it to occur; a key that only receives in-order events can + * retain the stale value indefinitely. */ private[autocdc] def coalesceIgnoredNulls( decomposedDf: DataFrame): DataFrame = @@ -1744,6 +1745,10 @@ object Scd2BatchProcessor { * * The supplied contexts must be nonempty and cover `path`. Struct fields without a context are * passed through unchanged. + * + * A non-nullable field raises `NOT_NULL_ASSERT_VIOLATION` when coalescing cannot supply it a + * non-null value: for example, when no preceding row authored it, when a delete boundary reset + * its inheritance, or when its null parent struct is rebuilt because a sibling leaf inherits. */ private def constructCoalescedIgnoreNullColumn( path: Seq[String], @@ -1752,7 +1757,7 @@ object Scd2BatchProcessor { val reconstructed = field.dataType match { case struct: StructType => // If this field is a struct, recursively reconstruct all of its children by applying their - // resolved values after ignore-null coalescing. + // resolved values after ignore-null coalescing. val contextsByChildName = contextsBeneath.groupBy(_.path(path.length)) val rebuilt = F.struct( struct.fields.toImmutableArraySeq.map { childField => @@ -1773,7 +1778,7 @@ object Scd2BatchProcessor { case _ => // If this field is not a struct (and therefore must be a leaf), either directly apply the // resolved value to inherit if the leaf should be inheriting, otherwise pass its value - // through as-is. + // through as-is. val context = contextsBeneath.head F.when(context.inherits, context.valueToInherit) .otherwise(F.col(QuotingUtils.quoteNameParts(path))) @@ -1784,7 +1789,7 @@ object Scd2BatchProcessor { } else { // If the field was marked as non-nullable but coalescing deduces it will resolve to a // null, throw an explicit exception. Pushing `AssertNotNull` into the plan for a - // non-nullable field also prevents Spark from preemptively complaining during analysis. + // non-nullable field also prevents Spark from preemptively complaining during analysis. ExpressionUtils.column( AssertNotNull(ExpressionUtils.expression(reconstructed), path)) } @@ -2225,7 +2230,7 @@ private[autocdc] object LeafInheritanceContext { // If this first row is an existing upsert, it cannot itself be re-coalesced in this sweep: // its predecessor is outside the affected window, so its stored value is the only safe // carry-in. With an unchanged selection that value is already correct. After a selection - // change, correcting the anchor is deferred until reconciliation includes preceding + // change, the anchor can only be corrected by a reconciliation that includes preceding // history. See the "Eventual consistency of the carry-in anchor" section in the // [[coalesceIgnoredNulls]] scaladoc. val rowEstablishesCarryIn = rowInheritanceContext.isFirstUpsertRepresentingRow diff --git a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala index 38678502551a7..8cd1719854e9e 100644 --- a/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -17,6 +17,7 @@ package org.apache.spark.sql.pipelines.autocdc +import org.apache.spark.SparkRuntimeException import org.apache.spark.sql.{functions => F, AnalysisException, QueryTest, Row} import org.apache.spark.sql.catalyst.util.QuotingUtils import org.apache.spark.sql.classic.DataFrame @@ -293,14 +294,8 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { StructField("id", IntegerType, nullable = false), StructField("profile", profileType, nullable = false, profileMetadata) )) - val authoredMap = versionMap( - Seq("profile", "address", "city") -> true, - Seq("profile", "address", "zip") -> true, - Seq("profile", "note") -> true) - val unauthoredZipMap = versionMap( - Seq("profile", "address", "city") -> true, - Seq("profile", "address", "zip") -> false, - Seq("profile", "note") -> true) + val authoredMap = versionMap() + val unauthoredZipMap = versionMap(Seq("profile", "address", "zip") -> false) val input = targetTableOf(schema)( Row(1, Row(Row("city-1", "zip-1"), "note-1"), 10L, null, cdcMetadata(10L, authoredMap)), @@ -322,6 +317,35 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { ) } + test("inheriting into a stored null struct fails when a non-nullable field is authored null") { + // Seq 20 stores profile = null, and its version map records profile.city as an authored null. + // profile.note was added by schema evolution, so it has no entry and inherits "note-1" from + // seq 10. Holding that inherited value requires a non-null profile struct, but city's + // authored null cannot be stored in a non-nullable field, so coalescing fails. This represents + // a schema misconfiguration issue, where the user is choosing a schema that simply cannot + // represent the coalesced result for this set and order of change events. This test locks in + // the error contract. + val profileType = new StructType() + .add("city", StringType, nullable = false) + .add("note", StringType) + val schema = new StructType() + .add("id", IntegerType) + .add("profile", profileType) + val input = targetTableOf(schema)( + Row(1, Row("city-1", "note-1"), 10L, null, cdcMetadata(10L, versionMap())), + Row(1, null, 20L, null, + cdcMetadata(20L, versionMap(Seq("profile", "city") -> true))) + ) + + checkError( + exception = intercept[SparkRuntimeException] { + coalesce(input, includeColumns("profile")).collect() + }, + condition = "NOT_NULL_ASSERT_VIOLATION", + parameters = Map("walkedTypePath" -> "\nprofile\ncity\n") + ) + } + gridTest("column selection honors the configured resolver")( Seq( (false, "VALUE"), @@ -602,12 +626,12 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { ) } - test("coalescing is eventually consistent for existing rows when " + - "ignore-null selection changes") { + test("a wider affected window corrects stale carry-in after selection changes") { // Demonstrates the "Eventual consistency of the carry-in anchor" contract of // [[Scd2BatchProcessor.coalesceIgnoredNulls]]: a leaf temporarily dropped from the ignore-null - // selection is not reconciled, so a later narrow affected window propagates a stale carry-in - // until a wider window pulls in the authoring row. + // selection is not reconciled, so a later narrow affected window propagates a stale carry-in. + // A window that includes the authoring row corrects it; this test constructs that window + // directly rather than proving one will occur. // "other" is authored and non-null in every row, so coalescing never changes it. It exists // only so the selection can move off "value" without becoming an (invalid) empty list. @@ -677,9 +701,9 @@ class Scd2CoalesceIgnoredNullsSuite extends QueryTest with SharedSparkSession { ) // -- Correction ------------------------------------------------------------------- - // A future batch's affected window expands to include seq 2. Now seq 2 is the - // first authored row in the window: its "new" propagates forward and the stale - // carry-in is corrected. + // Suppose a later batch's affected window includes seq 2. Now seq 2 is the first + // authored row in the window: its "new" propagates forward and the stale carry-in is + // corrected. val correctionInput = targetTableOf(schema)( Row(1, "new", "x", 2L, 3L, cdcMetadata(2L, emptyMap)), Row(1, "old", "x", 3L, 4L, cdcMetadata(3L, unauthoredMap)), From 79c971c75e1bc8aa05707fc31dc95085d9a7d9fe Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 17:00:41 +0000 Subject: [PATCH 12/13] format --- .../apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 05c6ee52353f8..5f73428cc9060 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -1016,7 +1016,7 @@ case class Scd2BatchProcessor( * interval gaps and delete boundaries that terminate inheritance, and before start/end * reconciliation so inherited tracked-history values determine the final SCD2 runs. * - * ==Eventual consistency of the carry-in anchor== + * ==== Eventual consistency of the carry-in anchor ==== * * A leaf's reconciliation in any pass is a function of two independent inputs: * From 4ea9a7d41abd746e2e26ad71413a27c796ab4035 Mon Sep 17 00:00:00 2001 From: Anish Mahto Date: Tue, 29 Sep 2026 19:17:41 +0000 Subject: [PATCH 13/13] add comments for updateVersionMapWithSchemaEvolution --- .../autocdc/Scd2BatchProcessor.scala | 21 ++++++++++++++----- 1 file changed, 16 insertions(+), 5 deletions(-) diff --git a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala index 5f73428cc9060..f7e0763643086 100644 --- a/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala +++ b/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd2BatchProcessor.scala @@ -1725,17 +1725,24 @@ object Scd2BatchProcessor { return versionMap } - val conditionalEntries = leafInheritanceContexts.map { context => + // Build the version map entry only for eligible schema evolved leaves, null for all other + // leaves. + val newEntryPerLeafOrNull = leafInheritanceContexts.map { context => F.when( context.needsSchemaEvolutionEntry, Scd2VersionMap.buildVersionMapEntry(context.path, authored = false)) } - val nonNullEntries = F.filter( - F.array(conditionalEntries: _*), (entry: Column) => entry.isNotNull) - val newEntries = F.map_from_entries(nonNullEntries) + + // Filter out the null entries, which explicitly represent leaves that don't need to gain a + // version map entry due to schema evolution. + val newEntries = F.filter( + F.array(newEntryPerLeafOrNull: _*), (entry: Column) => entry.isNotNull) + + // Concat new entries mapping with existing version map. F.when( versionMap.isNotNull, - F.map_concat(versionMap, newEntries)) + F.map_concat(versionMap, F.map_from_entries(newEntries)) + ) } /** @@ -2186,6 +2193,10 @@ private[autocdc] case class LeafInheritanceContext( inherits: Column, needsSchemaEvolutionEntry: Column) { + /** + * Name of the temporary column that holds [[valueToInheritIfAny]], the wrapped value this row + * may inherit, during ignore-null coalescing. + */ val valueToInheritIfAnyColName: String = LeafInheritanceContext.valueToInheritIfAnyColumnName(index)