AnishMahto commented on code in PR #56122:
URL: https://github.com/apache/spark/pull/56122#discussion_r3305529417
##########
sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/autocdc/Scd1BatchProcessor.scala:
##########
@@ -367,19 +367,29 @@ case class Scd1BatchProcessor(
val incomingWinsDelete = microbatchDeleteVersionField.isNotNull &&
microbatchDeleteVersionField > destinationUpsertVersionField
- // When the incoming upsert wins against an existing record, the entire
row (all columns)
- // will be overwritten, including the CDC metadata column. We only exclude
keys because
- // most merge implementations require that join columns are not being
mutated, even if
- // the mutation is a no-op.
val resolver = microbatchDf.sparkSession.sessionState.conf.resolver
val keyNames = changeArgs.keys.map(_.name)
+
+ def constructTargetColumnAssignmentsFromMicrobatch(columnName: String):
(String, Column) = {
+ // Map a column in the target table to its direct equivalent in the
microbatch. Note that
+ // because of target-table schema evolution during SDP dataset
materialization, the
+ // microbatch's columns are always a subset of (or equal to) the
target's columns.
+ val quotedCol = QuotingUtils.quoteIdentifier(columnName)
+ s"$destinationTableStr.$quotedCol" -> F.col(s"microbatch.$quotedCol")
+ }
+
+ // Most merge implementations require that join columns are not mutated,
even when the
+ // mutation would be a no-op. The remaining microbatch columns (including
the CDC metadata
+ // column) are overwritten outright when the incoming upsert wins.
val columnsToUpdateWhenIncomingWinsUpsert: Map[String, Column] =
microbatchDf.columns
.filterNot(c => keyNames.exists(resolver(_, c)))
- .map { c =>
- val quotedCol = QuotingUtils.quoteIdentifier(c)
- s"$destinationTableStr.$quotedCol" -> F.col(s"microbatch.$quotedCol")
- }
+ .map(constructTargetColumnAssignmentsFromMicrobatch)
+ .toMap
+
+ val columnsToInsertOnNewKey: Map[String, Column] =
+ microbatchDf.columns
+ .map(constructTargetColumnAssignmentsFromMicrobatch)
.toMap
Review Comment:
These changes were needed to support schema evolution, which I could only
test now that we've integrated flow execution with the rest of SDP.
Over multiple pipeline executions, the source microbatch's schema could
ultimately become a subset of the target table's schema - we should take care
to construct the column mappings appropriately.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]