Copilot commented on code in PR #13134:
URL: https://github.com/apache/gluten/pull/13134#discussion_r4109584954
##########
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:
`setNextObserver` delegates to `underlying`, but this decorator does not
override the corresponding `getNextObserver`, so the inherited chain field on
the wrapper remains empty and callers observe `None` even after setting a next
observer. That breaks the `ChainableExecutionObserver` contract and can also
hide the next wrapped observer from code that inspects the transaction's
observer; delegate the getter as well, wrapping the returned observer so native
transaction creation is retained.
--
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]