This is an automated email from the ASF dual-hosted git repository.
yew1eb pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/auron.git
The following commit(s) were added to refs/heads/master by this push:
new 572ce5a43 [AURON #2321] Support Iceberg column rename and
drop-then-add in the native scan (#2322)
572ce5a43 is described below
commit 572ce5a43a39175eeeb64a7e8d6cc6176662a21b
Author: linfeng <[email protected]>
AuthorDate: Wed Jul 1 14:47:28 2026 +0800
[AURON #2321] Support Iceberg column rename and drop-then-add in the native
scan (#2322)
# Which issue does this PR close?
Closes #2321
# Rationale for this change
The native Iceberg scan matches data-file columns by name, but Iceberg
tracks them by field-id. After a column rename, old files read as
all-NULL; after a drop-then-add of the same name, the new column reads
the old column's data.
# What changes are included in this PR?
Resolve columns by Iceberg field-id instead of by name:
- proto: add `field_id` to `Field`.
- JVM (`AuronIcebergSourceUtil`, `IcebergScanSupport`,
`NativeConverters`): extract top-level `name → field-id` from the scan's
`expectedSchema()` and serialize it into the plan.
- native (`auron-planner`, `scan/mod.rs`): stamp the id into Arrow field
metadata (`PARQUET:field_id`); `fields_match` matches by id when
present, else falls back to case-insensitive name matching (non-Iceberg
scans unchanged).
Nested-struct evolution and ORC rename/drop fall back to Spark, additive
evolution stays native.
# Are there any user-facing changes?
Yes. Iceberg queries on renamed or drop-then-added columns now return
correct results under the native scan. Unsupported cases fall back to
Spark. No API change.
# How was this patch tested?
Added cases to `AuronIcebergIntegrationSuite`
---
native-engine/auron-planner/proto/auron.proto | 2 +
native-engine/auron-planner/src/lib.rs | 51 +++++----
native-engine/datafusion-ext-plans/src/scan/mod.rs | 16 ++-
.../apache/spark/sql/auron/NativeConverters.scala | 14 ++-
.../execution/auron/plan/NativeGenerateBase.scala | 2 +-
.../spark/source/AuronIcebergSourceUtil.scala | 59 +++++++++++
.../sql/auron/iceberg/IcebergScanSupport.scala | 59 +++++++++--
.../auron/plan/NativeIcebergTableScanExec.scala | 3 +-
.../iceberg/AuronIcebergIntegrationSuite.scala | 118 +++++++++++++++++++++
9 files changed, 285 insertions(+), 39 deletions(-)
diff --git a/native-engine/auron-planner/proto/auron.proto
b/native-engine/auron-planner/proto/auron.proto
index 4a938f8ab..a905c8a36 100644
--- a/native-engine/auron-planner/proto/auron.proto
+++ b/native-engine/auron-planner/proto/auron.proto
@@ -826,6 +826,8 @@ message Field {
bool nullable = 3;
// for complex data types like structs, unions
repeated Field children = 4;
+ // Iceberg/Parquet field id. Zero means unset.
+ int32 field_id = 5;
}
message FixedSizeBinary {
diff --git a/native-engine/auron-planner/src/lib.rs
b/native-engine/auron-planner/src/lib.rs
index a0f7b83d2..c118cd2b0 100644
--- a/native-engine/auron-planner/src/lib.rs
+++ b/native-engine/auron-planner/src/lib.rs
@@ -13,10 +13,13 @@
// See the License for the specific language governing permissions and
// limitations under the License.
-use std::sync::Arc;
+use std::{collections::HashMap, sync::Arc};
use arrow::datatypes::{DataType, Field, Fields, IntervalUnit, Schema,
TimeUnit};
-use datafusion::{common::JoinSide, logical_expr::Operator,
scalar::ScalarValue};
+use datafusion::{
+ common::JoinSide, logical_expr::Operator,
parquet::arrow::PARQUET_FIELD_ID_META_KEY,
+ scalar::ScalarValue,
+};
use datafusion_ext_plans::{agg::AggFunction, joins::join_utils::JoinType};
use crate::error::PlanSerDeError;
@@ -406,17 +409,29 @@ impl TryInto<DataType> for &Box<protobuf::List> {
impl TryInto<Field> for &protobuf::Field {
type Error = PlanSerDeError;
fn try_into(self) -> Result<Field, Self::Error> {
- let pb_datatype = self.arrow_type.as_ref().ok_or_else(|| {
- proto_error(
- "Protobuf deserialization error: Field message missing
required field 'arrow_type'",
- )
- })?;
+ build_arrow_field(self)
+ }
+}
+
+fn build_arrow_field(field: &protobuf::Field) -> Result<Field, PlanSerDeError>
{
+ let pb_datatype = field.arrow_type.as_ref().ok_or_else(|| {
+ proto_error(
+ "Protobuf deserialization error: Field message missing required
field 'arrow_type'",
+ )
+ })?;
+ let arrow_field = Field::new(
+ field.name.as_str(),
+ pb_datatype.as_ref().try_into()?,
+ field.nullable,
+ );
- Ok(Field::new(
- self.name.as_str(),
- pb_datatype.as_ref().try_into()?,
- self.nullable,
- ))
+ if field.field_id == 0 {
+ Ok(arrow_field)
+ } else {
+ Ok(arrow_field.with_metadata(HashMap::from([(
+ PARQUET_FIELD_ID_META_KEY.to_string(),
+ field.field_id.to_string(),
+ )])))
}
}
@@ -427,17 +442,7 @@ impl TryInto<Schema> for &protobuf::Schema {
let fields = self
.columns
.iter()
- .map(|c| {
- let pb_arrow_type_res = c
- .arrow_type
- .as_ref()
- .ok_or_else(|| proto_error("Protobuf deserialization
error: Field message was missing required field 'arrow_type'"));
- let pb_arrow_type: &protobuf::ArrowType = match
pb_arrow_type_res {
- Ok(res) => res,
- Err(e) => return Err(e),
- };
- Ok(Field::new(&c.name, pb_arrow_type.try_into()?, c.nullable))
- })
+ .map(build_arrow_field)
.collect::<Result<Vec<_>, _>>()?;
Ok(Schema::new(fields))
}
diff --git a/native-engine/datafusion-ext-plans/src/scan/mod.rs
b/native-engine/datafusion-ext-plans/src/scan/mod.rs
index 3d7899acd..257999980 100644
--- a/native-engine/datafusion-ext-plans/src/scan/mod.rs
+++ b/native-engine/datafusion-ext-plans/src/scan/mod.rs
@@ -25,6 +25,7 @@ use datafusion::{
datasource::schema_adapter::{
SchemaAdapter, SchemaAdapterFactory, SchemaMapper, SchemaMapping,
},
+ parquet::arrow::PARQUET_FIELD_ID_META_KEY,
};
use datafusion_ext_commons::df_execution_err;
@@ -57,11 +58,10 @@ impl SchemaAdapter for AuronSchemaAdapter {
fn map_column_index(&self, index: usize, file_schema: &Schema) ->
Option<usize> {
let field = self.table_schema.field(index);
- // use case insensitive matching
file_schema
.fields()
.iter()
- .position(|f| f.name().eq_ignore_ascii_case(field.name()))
+ .position(|file_field| fields_match(field, file_field))
}
fn map_schema(&self, file_schema: &Schema) -> Result<(Arc<dyn
SchemaMapper>, Vec<usize>)> {
@@ -73,7 +73,7 @@ impl SchemaAdapter for AuronSchemaAdapter {
.table_schema
.fields()
.iter()
- .position(|f| f.name().eq_ignore_ascii_case(file_field.name()))
+ .position(|table_field| fields_match(table_field, file_field))
{
field_mappings[table_idx] = Some(projection.len());
projection.push(file_idx);
@@ -89,6 +89,16 @@ impl SchemaAdapter for AuronSchemaAdapter {
}
}
+fn fields_match(table_field: &Field, file_field: &Field) -> bool {
+ match table_field.metadata().get(PARQUET_FIELD_ID_META_KEY) {
+ Some(table_field_id) => match
file_field.metadata().get(PARQUET_FIELD_ID_META_KEY) {
+ Some(file_field_id) => file_field_id == table_field_id,
+ None => table_field.name().eq_ignore_ascii_case(file_field.name()),
+ },
+ None => table_field.name().eq_ignore_ascii_case(file_field.name()),
+ }
+}
+
pub fn create_auron_schema_mapper(
table_schema: &SchemaRef,
field_mappings: &[Option<usize>],
diff --git
a/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala
b/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala
index 5dd768309..378a8d662 100644
---
a/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala
+++
b/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala
@@ -216,18 +216,22 @@ object NativeConverters extends Logging {
arrowTypeBuilder.build()
}
- def convertField(sparkField: StructField): pb.Field = {
- pb.Field
+ def convertField(sparkField: StructField, fieldId: Option[Int] = None):
pb.Field = {
+ val fieldBuilder = pb.Field
.newBuilder()
.setName(sparkField.name)
.setNullable(sparkField.nullable)
.setArrowType(convertDataType(sparkField.dataType))
- .build()
+ fieldId.foreach(fieldBuilder.setFieldId)
+ fieldBuilder.build()
}
- def convertSchema(sparkSchema: StructType): pb.Schema = {
+ def convertSchema(
+ sparkSchema: StructType,
+ fieldIdsByName: Map[String, Int] = Map.empty): pb.Schema = {
val schemaBuilder = pb.Schema.newBuilder()
- sparkSchema.foreach(sparkField =>
schemaBuilder.addColumns(convertField(sparkField)))
+ sparkSchema.foreach(sparkField =>
+ schemaBuilder.addColumns(convertField(sparkField,
fieldIdsByName.get(sparkField.name))))
schemaBuilder.build()
}
diff --git
a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateBase.scala
b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateBase.scala
index a607b2894..202645b6a 100644
---
a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateBase.scala
+++
b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateBase.scala
@@ -118,7 +118,7 @@ abstract class NativeGenerateBase(
}
private def nativeGeneratorOutput =
- Util.getSchema(generatorOutput).map(NativeConverters.convertField)
+ Util.getSchema(generatorOutput).map(field =>
NativeConverters.convertField(field))
private def nativeRequiredChildOutput =
Util.getSchema(requiredChildOutput).map(_.name)
diff --git
a/thirdparty/auron-iceberg/src/main/scala/org/apache/iceberg/spark/source/AuronIcebergSourceUtil.scala
b/thirdparty/auron-iceberg/src/main/scala/org/apache/iceberg/spark/source/AuronIcebergSourceUtil.scala
index 6b04f16a3..c1d0a58d6 100644
---
a/thirdparty/auron-iceberg/src/main/scala/org/apache/iceberg/spark/source/AuronIcebergSourceUtil.scala
+++
b/thirdparty/auron-iceberg/src/main/scala/org/apache/iceberg/spark/source/AuronIcebergSourceUtil.scala
@@ -16,8 +16,14 @@
*/
package org.apache.iceberg.spark.source
+import scala.collection.JavaConverters._
+
+import org.apache.iceberg.types.TypeUtil
+
object AuronIcebergSourceUtil {
+ final case class RenameOrDrop(topLevel: Boolean, nested: Boolean)
+
def getClassOfSparkBatchQueryScan(): Class[SparkBatchQueryScan] = {
classOf[SparkBatchQueryScan]
}
@@ -25,4 +31,57 @@ object AuronIcebergSourceUtil {
def getClassOfSparkInputPartition(): Class[SparkInputPartition] = {
classOf[SparkInputPartition]
}
+
+ def expectedFieldIds(scan: AnyRef): Map[String, Int] = {
+ val expectedSchema = asBatchQueryScan(scan).expectedSchema()
+ expectedSchema.columns().asScala.map(field => field.name() ->
field.fieldId()).toMap
+ }
+
+ def detectRenameOrDrop(scan: AnyRef): RenameOrDrop = {
+ val table = asBatchQueryScan(scan).table()
+ val currentFields = collectFieldIdToName(table.schema())
+
+ table
+ .schemas()
+ .asScala
+ .filterNot(_._1 == table.schema().schemaId())
+ .values
+ .foldLeft(RenameOrDrop(topLevel = false, nested = false)) { (result,
schema) =>
+ collectFieldIdToName(schema).foldLeft(result) {
+ case (currentResult, (fieldId, historicalField)) =>
+ currentFields.get(fieldId) match {
+ case Some(currentField) if currentField.name !=
historicalField.name =>
+ if (historicalField.topLevel || currentField.topLevel) {
+ currentResult.copy(topLevel = true)
+ } else {
+ currentResult.copy(nested = true)
+ }
+ case None =>
+ if (historicalField.topLevel) {
+ currentResult.copy(topLevel = true)
+ } else {
+ currentResult.copy(nested = true)
+ }
+ case _ =>
+ currentResult
+ }
+ }
+ }
+ }
+
+ final private case class FieldIdentity(name: String, topLevel: Boolean)
+
+ private def collectFieldIdToName(schema: org.apache.iceberg.Schema):
Map[Int, FieldIdentity] = {
+ val topLevelFieldIds = schema.columns().asScala.map(_.fieldId()).toSet
+ TypeUtil
+ .indexById(schema.asStruct())
+ .asScala
+ .map { case (fieldId, field) =>
+ fieldId.toInt -> FieldIdentity(field.name(),
topLevelFieldIds.contains(fieldId.toInt))
+ }
+ .toMap
+ }
+
+ private def asBatchQueryScan(scan: AnyRef): SparkBatchQueryScan =
+ scan.asInstanceOf[SparkBatchQueryScan]
}
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 d50cd945b..3aa85b2de 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
@@ -25,6 +25,7 @@ import org.apache.iceberg.spark.source.AuronIcebergSourceUtil
import org.apache.spark.internal.Logging
import org.apache.spark.sql.auron.NativeConverters
import org.apache.spark.sql.catalyst.expressions.{And => SparkAnd,
AttributeReference, EqualTo, Expression => SparkExpression, GreaterThan,
GreaterThanOrEqual, In, IsNaN, IsNotNull, IsNull, LessThan, LessThanOrEqual,
Literal, Not => SparkNot, Or => SparkOr}
+import org.apache.spark.sql.catalyst.trees.TreeNodeTag
import org.apache.spark.sql.connector.read.{InputPartition, Scan}
import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
import org.apache.spark.sql.internal.SQLConf
@@ -48,9 +49,12 @@ final case class IcebergScanPlan(
readSchema: StructType,
fileSchema: StructType,
partitionSchema: StructType,
- pruningPredicates: Seq[pb.PhysicalExprNode])
+ pruningPredicates: Seq[pb.PhysicalExprNode],
+ fieldIdsByName: Map[String, Int])
object IcebergScanSupport extends Logging {
+ private val scanPlanTag: TreeNodeTag[Option[IcebergScanPlan]] = TreeNodeTag(
+ "auron.iceberg.scan.plan")
private val SparkChangelogScanClassName =
"org.apache.iceberg.spark.source.SparkChangelogScan"
@@ -79,6 +83,16 @@ object IcebergScanSupport extends Logging {
}
def plan(exec: BatchScanExec): Option[IcebergScanPlan] = {
+ exec.getTagValue(scanPlanTag) match {
+ case Some(cached) => cached
+ case None =>
+ val planned = planUncached(exec)
+ exec.setTagValue(scanPlanTag, planned)
+ planned
+ }
+ }
+
+ private def planUncached(exec: BatchScanExec): Option[IcebergScanPlan] = {
val scan = exec.scan
val scanClassName = scan.getClass.getName
// Only handle Iceberg scans; other sources must stay on Spark's path.
@@ -104,6 +118,31 @@ object IcebergScanSupport extends Logging {
}
val (fileSchema, partitionSchema) = schemas.get
+ val fieldIdsByName =
+ try {
+ AuronIcebergSourceUtil.expectedFieldIds(scan.asInstanceOf[AnyRef])
+ } catch {
+ case NonFatal(t) =>
+ logWarning(s"Failed to inspect Iceberg field ids for
$scanClassName.", t)
+ return None
+ }
+
+ val renameOrDrop =
+ try {
+ AuronIcebergSourceUtil.detectRenameOrDrop(scan.asInstanceOf[AnyRef])
+ } catch {
+ case NonFatal(t) =>
+ logWarning(s"Failed to inspect Iceberg schema history for
$scanClassName.", t)
+ return None
+ }
+ assert(!renameOrDrop.nested, "Nested Iceberg rename or drop is not
supported.")
+
+ val missingFieldIds =
+ fileSchema.fields.filterNot(field =>
fieldIdsByName.contains(field.name)).map(_.name)
+ assert(
+ missingFieldIds.isEmpty,
+ s"Missing Iceberg field ids for columns: ${missingFieldIds.mkString(",
")}")
+
val partitions = inputPartitions(exec)
// Empty scan (e.g. empty table) should still build a plan to return no
rows.
if (partitions.isEmpty) {
@@ -115,7 +154,8 @@ object IcebergScanSupport extends Logging {
readSchema,
fileSchema,
partitionSchema,
- Seq.empty))
+ Seq.empty,
+ fieldIdsByName))
}
val icebergPartitions = partitions.flatMap(icebergPartition)
@@ -142,7 +182,11 @@ object IcebergScanSupport extends Logging {
}
val format = formats.headOption.getOrElse(FileFormat.PARQUET)
- if (format != FileFormat.PARQUET && format != FileFormat.ORC) {
+ // ORC cannot match Iceberg columns by field-id yet, so any historical
top-level
+ // rename/drop may make older ORC files unsafe for native name/position
matching.
+ val supportedFormat =
+ format == FileFormat.PARQUET || (format == FileFormat.ORC &&
!renameOrDrop.topLevel)
+ if (!supportedFormat) {
return None
}
@@ -155,7 +199,8 @@ object IcebergScanSupport extends Logging {
readSchema,
fileSchema,
partitionSchema,
- pruningPredicates))
+ pruningPredicates,
+ fieldIdsByName))
}
private def planChangelogScan(exec: BatchScanExec, scan: Scan):
Option[IcebergScanPlan] = {
@@ -175,7 +220,8 @@ object IcebergScanSupport extends Logging {
readSchema,
fileSchema,
partitionSchema,
- Seq.empty))
+ Seq.empty,
+ Map.empty))
}
val icebergPartitions = partitions.flatMap(icebergPartition)
@@ -223,7 +269,8 @@ object IcebergScanSupport extends Logging {
readSchema,
fileSchema,
partitionSchema,
- pruningPredicates))
+ pruningPredicates,
+ Map.empty))
}
private def supportedSchemas(
diff --git
a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeIcebergTableScanExec.scala
b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeIcebergTableScanExec.scala
index 831a9ba21..3dfa08b65 100644
---
a/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeIcebergTableScanExec.scala
+++
b/thirdparty/auron-iceberg/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeIcebergTableScanExec.scala
@@ -74,7 +74,8 @@ case class NativeIcebergTableScanExec(basedScan:
BatchScanExec, plan: IcebergSca
}
private lazy val fileSizes: Map[String, Long] = buildFileSizes()
- private lazy val nativeFileSchema: pb.Schema =
NativeConverters.convertSchema(fileSchema)
+ private lazy val nativeFileSchema: pb.Schema =
+ NativeConverters.convertSchema(fileSchema, plan.fieldIdsByName)
private lazy val nativePartitionSchema: pb.Schema =
NativeConverters.convertSchema(partitionSchema)
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 2412ba8a4..1142b6e03 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
@@ -190,6 +190,124 @@ class AuronIcebergIntegrationSuite
}
}
+ test("iceberg native parquet scan reads top-level renamed columns by field
id") {
+ withTable("local.db.t_rename") {
+ sql("create table local.db.t_rename (id int, old_name string) using
iceberg")
+ sql("insert into local.db.t_rename values (1, 'before')")
+ sql("alter table local.db.t_rename rename column old_name to new_name")
+ sql("insert into local.db.t_rename values (2, 'after')")
+
+ val df = sql("select id, new_name from local.db.t_rename")
+ checkAnswer(df, Seq(Row(1, "before"), Row(2, "after")))
+
assert(df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg native parquet scan does not reuse a dropped field id for an
added column") {
+ withTable("local.db.t_drop_add") {
+ sql("create table local.db.t_drop_add (id int, value string) using
iceberg")
+ sql("insert into local.db.t_drop_add values (1, 'old')")
+ sql("alter table local.db.t_drop_add drop column value")
+ sql("alter table local.db.t_drop_add add column value string")
+ sql("insert into local.db.t_drop_add values (2, 'new')")
+
+ val df = sql("select id, value from local.db.t_drop_add")
+ checkAnswer(df, Seq(Row(1, null), Row(2, "new")))
+
assert(df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg ORC scan falls back after a top-level column rename") {
+ withTable("local.db.t_orc_rename") {
+ sql("""
+ |create table local.db.t_orc_rename (id int, old_name string)
+ |using iceberg
+ |tblproperties ('write.format.default' = 'orc')
+ |""".stripMargin)
+ sql("insert into local.db.t_orc_rename values (1, 'before')")
+ sql("alter table local.db.t_orc_rename rename column old_name to
new_name")
+
+ val df = sql("select id, new_name from local.db.t_orc_rename")
+ checkAnswer(df, Seq(Row(1, "before")))
+
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg ORC scan falls back after top-level drop and add with the same
name") {
+ withTable("local.db.t_orc_drop_add") {
+ sql("""
+ |create table local.db.t_orc_drop_add (id int, value string)
+ |using iceberg
+ |tblproperties ('write.format.default' = 'orc')
+ |""".stripMargin)
+ sql("insert into local.db.t_orc_drop_add values (1, 'old')")
+ sql("alter table local.db.t_orc_drop_add drop column value")
+ sql("alter table local.db.t_orc_drop_add add column value string")
+ sql("insert into local.db.t_orc_drop_add values (2, 'new')")
+
+ val df = sql("select id, value from local.db.t_orc_drop_add")
+ checkAnswer(df, Seq(Row(1, null), Row(2, "new")))
+
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg ORC scan remains native for additive schema evolution") {
+ withTable("local.db.t_orc_add") {
+ sql("""
+ |create table local.db.t_orc_add (id int)
+ |using iceberg
+ |tblproperties ('write.format.default' = 'orc')
+ |""".stripMargin)
+ sql("insert into local.db.t_orc_add values (1)")
+ sql("alter table local.db.t_orc_add add column value string")
+
+ val df = sql("select id, value from local.db.t_orc_add")
+ checkAnswer(df, Seq(Row(1, null)))
+
assert(df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg scan falls back after a nested column rename") {
+ withTable("local.db.t_nested_rename") {
+ sql("""
+ |create table local.db.t_nested_rename (
+ | id int,
+ | payload struct<old_name:string>
+ |) using iceberg
+ |""".stripMargin)
+ sql("insert into local.db.t_nested_rename values (1,
named_struct('old_name', 'before'))")
+ sql("alter table local.db.t_nested_rename rename column payload.old_name
to new_name")
+
+ val df = sql("select id, payload.new_name from local.db.t_nested_rename")
+ checkAnswer(df, Seq(Row(1, "before")))
+
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
+ test("iceberg scan falls back when top-level and nested columns are both
renamed") {
+ withTable("local.db.t_top_and_nested_rename") {
+ sql("""
+ |create table local.db.t_top_and_nested_rename (
+ | old_id int,
+ | payload struct<old_name:string>
+ |) using iceberg
+ |""".stripMargin)
+ sql("""insert into local.db.t_top_and_nested_rename
+ |values (1, named_struct('old_name', 'before'))
+ |""".stripMargin)
+ sql("alter table local.db.t_top_and_nested_rename rename column old_id
to new_id")
+ sql("""
+ |alter table local.db.t_top_and_nested_rename
+ |rename column payload.old_name to new_name
+ |""".stripMargin)
+
+ val df =
+ sql("select new_id, payload.new_name from
local.db.t_top_and_nested_rename")
+ checkAnswer(df, Seq(Row(1, "before")))
+
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
+ }
+ }
+
test("iceberg native scan is applied when delete files are null (format
v1)") {
withTable("local.db.t_v1") {
sql("""