felipepessoto commented on code in PR #13134:
URL: https://github.com/apache/gluten/pull/13134#discussion_r4109626901
##########
backends-velox/src-delta33/main/scala/org/apache/spark/sql/execution/datasources/v2/DeltaWriteOperators.scala:
##########
@@ -77,26 +74,79 @@ case class GlutenDeltaRunnableCommand(delegate:
RunnableCommand) extends LeafRun
}
object DeltaV2WriteOperators {
- object UseColumnarDeltaTransactionLog extends TransactionExecutionObserver {
+ private[sql] def withColumnarTransaction[T](f: => T): T = {
+ TransactionExecutionObserver.getObserver match {
+ case _: UseColumnarDeltaTransactionLog => f
+ case observer =>
+ TransactionExecutionObserver.setObserver(wrap(observer))
+ try {
+ f
+ } finally {
+ // A completed transaction may have advanced to the next observer.
+
TransactionExecutionObserver.setObserver(unwrap(TransactionExecutionObserver.getObserver))
+ }
+ }
+ }
+
+ private def wrap(observer: TransactionExecutionObserver):
TransactionExecutionObserver =
+ observer match {
+ case _: UseColumnarDeltaTransactionLog => observer
+ case _ => new UseColumnarDeltaTransactionLog(observer)
+ }
+
+ private def unwrap(observer: TransactionExecutionObserver):
TransactionExecutionObserver =
+ observer match {
+ case columnar: UseColumnarDeltaTransactionLog => columnar.underlying
+ case _ => observer
+ }
+
+ private class UseColumnarDeltaTransactionLog(val underlying:
TransactionExecutionObserver)
+ extends TransactionExecutionObserver {
override def startingTransaction(f: => OptimisticTransaction):
OptimisticTransaction = {
- val delegate = f
- new GlutenOptimisticTransaction(delegate)
+ underlying.startingTransaction {
+ new GlutenOptimisticTransaction(f)
+ }
}
- override def preparingCommit[T](f: => T): T = f
+ override def preparingCommit[T](f: => T): T = underlying.preparingCommit(f)
- override def beginDoCommit(): Unit = ()
+ override def beginDoCommit(): Unit = underlying.beginDoCommit()
- override def beginBackfill(): Unit = ()
+ override def beginBackfill(): Unit = underlying.beginBackfill()
- override def beginPostCommit(): Unit = ()
+ override def beginPostCommit(): Unit = underlying.beginPostCommit()
- override def transactionCommitted(): Unit = ()
+ override def transactionCommitted(): Unit = withObserverAdvance {
+ underlying.transactionCommitted()
+ }
- override def transactionAborted(): Unit = ()
+ override def transactionAborted(): Unit = withObserverAdvance {
+ underlying.transactionAborted()
+ }
override def createChild(): TransactionExecutionObserver = {
- TransactionExecutionObserver.getObserver
+ wrap(underlying.createChild())
+ }
+
+ override def setNextObserver(nextTxnObserver:
TransactionExecutionObserver): Unit = {
+ underlying.setNextObserver(unwrap(nextTxnObserver))
+ }
Review Comment:
Thanks for the review. I checked the compiled Delta 3.3.2 / 4.0.1 APIs and
the Delta 4.1.0 / 4.2.0 sources: these versions do not define a
`getNextObserver` method. [`nextObserver` is protected Scala
state](https://github.com/delta-io/delta/blob/v4.2.0/spark/src/main/scala/org/apache/spark/sql/delta/TransactionExecutionObserver.scala#L19-L34),
not a public getter in the Scala contract.
Both `setNextObserver` and `advanceToNextThreadObserver` delegate to the
same underlying observer. The phase-locking commit/abort callbacks also execute
on that underlying observer, so they read its chain state rather than the
unused inherited field on the decorator. When a callback advances the
thread-local observer, the decorator rewraps that replacement to retain native
transaction creation.
The [regression
coverage](https://github.com/apache/gluten/blob/ed6cd7bada3377751ae5fab47994144436a29cb9/backends-velox/src-delta33/test/scala/org/apache/spark/sql/delta/DeltaTransactionObserverSuite.scala#L143-L192)
exercises commit-driven and abort-driven advancement using the real Delta
phase observer, asserts that the next transaction is still native, and covers
explicit advancement through the decorator. These cases pass on both Delta 3.3
and 4.0. I do not think a getter override is applicable to the currently
supported API, so I am leaving this unchanged.
--
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]