Skip to content

Commit b7ce3c2

Browse files
committed
[AURON #2253] Support insert-only Iceberg changelog native scan
1 parent 161d19c commit b7ce3c2

6 files changed

Lines changed: 390 additions & 106 deletions

File tree

spark-extension/src/main/scala/org/apache/spark/sql/auron/AuronConvertStrategy.scala

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,9 @@ object AuronConvertStrategy extends Logging {
7171
if (!exec.getTagValue(neverConvertReasonTag).isDefined) {
7272
exec.setTagValue(
7373
neverConvertReasonTag,
74-
s"${exec.getClass.getSimpleName} is not supported yet.")
74+
converted
75+
.getTagValue(neverConvertReasonTag)
76+
.getOrElse(s"${exec.getClass.getSimpleName} is not supported yet."))
7577
}
7678
}
7779
danglingChildren = newDangling :+ converted

spark-extension/src/main/scala/org/apache/spark/sql/auron/AuronConverters.scala

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@ import org.apache.spark.sql.execution.auron.plan.NativeSortBase
6565
import org.apache.spark.sql.execution.auron.plan.NativeUnionBase
6666
import org.apache.spark.sql.execution.auron.plan.Util
6767
import org.apache.spark.sql.execution.command.DataWritingCommandExec
68+
import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
6869
import org.apache.spark.sql.execution.exchange.BroadcastExchangeExec
6970
import org.apache.spark.sql.execution.exchange.ShuffleExchangeExec
7071
import org.apache.spark.sql.execution.joins._
@@ -415,6 +416,8 @@ object AuronConverters extends Logging {
415416
} else {
416417
s"Falling back exec: ${exec.getClass.getSimpleName}: ${e.getMessage}"
417418
}
419+
case _: BatchScanExec =>
420+
s"${e.getMessage.replaceFirst("^assertion failed: ?", "")}"
418421
case _ =>
419422
s"Falling back exec: ${exec.getClass.getSimpleName}: ${e.getMessage}"
420423
}

thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergConvertProvider.scala

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,8 @@ class IcebergConvertProvider extends AuronConvertProvider with Logging {
4444

4545
override def isSupported(exec: SparkPlan): Boolean = {
4646
exec match {
47-
case e: BatchScanExec => IcebergScanSupport.plan(e).nonEmpty
47+
case e: BatchScanExec if IcebergScanSupport.isIcebergScan(e.scan) =>
48+
IcebergScanSupport.plan(e).nonEmpty || IcebergScanSupport.fallbackReason(e).nonEmpty
4849
case _ => false
4950
}
5051
}
@@ -56,7 +57,10 @@ class IcebergConvertProvider extends AuronConvertProvider with Logging {
5657
case Some(plan) =>
5758
AuronConverters.addRenameColumnsExec(NativeIcebergTableScanExec(e, plan))
5859
case None =>
59-
exec
60+
IcebergScanSupport.fallbackReason(e) match {
61+
case Some(reason) => throw new AssertionError(reason)
62+
case None => exec
63+
}
6064
}
6165
case _ => exec
6266
}

0 commit comments

Comments
 (0)