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..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 @@ -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} @@ -967,6 +968,193 @@ case class Scd2BatchProcessor( .drop(Scd2BatchProcessor.isRedundantDeleteEncodingColName) } + /** + * 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, + 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 nearest preceding authoring + * 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 + * 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. + * + * ==== 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, 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. + * + * 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 = + 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 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 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. * @@ -1529,6 +1717,92 @@ 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 + } + + // 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)) + } + + // 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, F.map_from_entries(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. + * + * 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], + field: StructField, + 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 => + val childPath = path :+ childField.name + contextsByChildName + .get(childField.name) + .map(constructCoalescedIgnoreNullColumn(childPath, childField, _)) + .getOrElse(F.col(QuotingUtils.quoteNameParts(childPath))) + .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 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)) + } + validatedReconstructed.cast(field.dataType) + } + /** * Name of temporary column projected onto microbatch to compute the min sequencing value per * key within the microbatch. @@ -1669,11 +1943,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 +2007,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 +2086,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 +2123,181 @@ 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) { + + /** + * 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) + + /** + * 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. + // + // 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, 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 + 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..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,7 +116,10 @@ case class Scd2ForeachBatchHandler( .transform(d => batchProcessor.assertWellFormedRowsPostDecomposition(d, batchId)) .transform(batchProcessor.dropRedundantRowsPostDecomposition) - val reconciledAndRoutedDf = decomposedDf + val withCoalescedIgnoredNullsDf = + batchProcessor.coalesceIgnoredNulls(decomposedDf) + + val reconciledAndRoutedDf = withCoalescedIgnoredNullsDf .transform(batchProcessor.reconcileStartAndEndAt) .transform(batchProcessor.dropLeftoverDeletesPostReconciliation) .transform(batchProcessor.promoteDecompositionTailsToTombstones) 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. */ 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..8cd1719854e9e --- /dev/null +++ b/sql/pipelines/src/test/scala/org/apache/spark/sql/pipelines/autocdc/Scd2CoalesceIgnoredNullsSuite.scala @@ -0,0 +1,721 @@ +/* + * 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.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 +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) + + 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)) + ) + + val firstPass = coalesce(input, selection) + checkAnswer( + 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") { + val selection = includeColumns("value") + val schema = new StructType() + .add("id", IntegerType) + .add("value", StringType) + 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, "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, "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)), + Row(3, "stale", 30L, null, cdcMetadata(30L, unauthoredMap)) + ) + + checkAnswer( + coalesce(input, selection), + Seq( + Row(1, "tombstone", 20L, 20L, cdcMetadata(20L, null)), + Row(1, null, 10L, null, cdcMetadata(30L, unauthoredMap)), + 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)) + ) + ) + } + + 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)) + ) + ) + } + + 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() + 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)), + 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)) + ) + ) + } + + 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"), + (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())) + ) + + 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" + ) + ) + } + } + + 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)) + ) + ) + } + + 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. + // 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. + 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 ------------------------------------------------------------------- + // 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)), + 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)) + ) + ) + } +} 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