From 6c61e738a0423caee91398990b37190e8cb082bd Mon Sep 17 00:00:00 2001 From: weimingdiit Date: Sat, 25 Jul 2026 11:19:35 +0800 Subject: [PATCH] [AURON #2427] Report fallback reason for unsupported Iceberg scan data types Signed-off-by: weimingdiit --- .../auron/iceberg/IcebergScanSupport.scala | 19 ++++++++++++++++++- .../AuronIcebergIntegrationSuite.scala | 6 ++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala index 31d877930..e280ab1e4 100644 --- a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala +++ b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/auron/iceberg/IcebergScanSupport.scala @@ -83,7 +83,14 @@ object IcebergScanSupport extends Logging { if (collectUnsupportedMetadataColumns(scan.readSchema, isChangelogScan).nonEmpty) { Some("Has per-row materialization (for example _pos).") } else { - None + val unsupportedFields = collectUnsupportedDataTypeFields(scan.readSchema, isChangelogScan) + if (unsupportedFields.nonEmpty) { + Some( + s"Unsupported Iceberg scan schema. Unsupported fields/types: " + + s"${unsupportedFields.mkString(", ")}.") + } else { + None + } } } @@ -377,6 +384,16 @@ object IcebergScanSupport extends Logging { field.name } + private def collectUnsupportedDataTypeFields( + schema: StructType, + isChangelogScan: Boolean): Seq[String] = + schema.fields + .filterNot(field => + isIcebergMetadataColumn(field.name, isChangelogScan) && + !isSupportedMetadataColumn(field, isChangelogScan)) + .filterNot(field => NativeConverters.isTypeSupported(field.dataType)) + .map(field => s"${field.name}: ${field.dataType.catalogString}") + private def isIcebergMetadataColumn(name: String, isChangelogScan: Boolean): Boolean = MetadataColumns.isMetadataColumn(name) || (isChangelogScan && ChangelogMetadataColumnNames.contains(name)) diff --git a/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala b/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala index 6d342bd48..1140b67b6 100644 --- a/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala +++ b/thirdparty/auron-iceberg/src/test/scala/org/apache/auron/iceberg/AuronIcebergIntegrationSuite.scala @@ -733,6 +733,12 @@ class AuronIcebergIntegrationSuite checkAnswer(df, Seq(Row(1, new java.math.BigDecimal("123.4500000000")))) val plan = df.queryExecution.executedPlan.toString() assert(!plan.contains("NativeIcebergTableScan")) + val neverConvertReasonTag: TreeNodeTag[String] = TreeNodeTag("auron.never.convert.reason") + assert( + collectFirst(df.queryExecution.executedPlan) { case batchScanExec: BatchScanExec => + batchScanExec.getTagValue(neverConvertReasonTag) + }.get.get.equals( + "Unsupported Iceberg scan schema. Unsupported fields/types: amount: decimal(38,10).")) } }