Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,8 +19,8 @@ package org.apache.iceberg.spark.source
import scala.collection.JavaConverters._

import org.apache.commons.lang3.reflect.FieldUtils
import org.apache.iceberg.Table
import org.apache.iceberg.types.TypeUtil
import org.apache.iceberg.{MetadataColumns, Partitioning, Table}
import org.apache.iceberg.types.{Types, TypeUtil}

object AuronIcebergSourceUtil {

Expand All @@ -39,6 +39,22 @@ object AuronIcebergSourceUtil {
expectedSchema.columns().asScala.map(field => field.name() -> field.fieldId()).toMap
}

def expectedPartitionType(scan: AnyRef): Option[Types.StructType] = {
val batchScan = asBatchQueryScan(scan)
Option(batchScan.expectedSchema().findType(MetadataColumns.PARTITION_COLUMN_ID)).map {
projected =>
// Spark may renumber nested metadata fields; retain Iceberg partition field IDs.
val partitionType = Partitioning.partitionType(batchScan.table())
Types.StructType.of(
projected
.asStructType()
.fields()
.asScala
.map(field => partitionType.field(field.name()))
.asJava)
}
}

def expectedFieldIdsForChangelogScan(scan: AnyRef): Map[String, Int] = {
// SparkChangelogScan does not expose Iceberg expectedSchema/table accessors.
// Keep the internal field-name assumptions localized here; callers fallback if they change.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import org.apache.iceberg.{AddedRowsScanTask, ChangelogOperation, ChangelogScanT
import org.apache.iceberg.expressions.{And => IcebergAnd, BoundPredicate, Expression => IcebergExpression, Not => IcebergNot, Or => IcebergOr, UnboundPredicate}
import org.apache.iceberg.spark.SparkUtil
import org.apache.iceberg.spark.source.AuronIcebergSourceUtil
import org.apache.iceberg.types.Types
import org.apache.iceberg.util.PartitionUtil
import org.apache.spark.internal.Logging
import org.apache.spark.sql.auron.{NativeConverters, Shims}
Expand Down Expand Up @@ -246,7 +247,9 @@ object IcebergScanSupport extends Logging {
}

val pruningPredicates = collectPruningPredicates(scan.asInstanceOf[AnyRef], readSchema)
val nativeTasks = fileTasks.map(task => toNativeScanTask(task, partitionSchema))
val partitionType = AuronIcebergSourceUtil.expectedPartitionType(scan.asInstanceOf[AnyRef])
val nativeTasks = fileTasks.map(task =>
toNativeScanTask(task, partitionSchema, fieldIdsByName, partitionType.orNull))
withIdentityPartitions(
IcebergScanPlan(
nativeTasks,
Expand Down Expand Up @@ -436,6 +439,7 @@ object IcebergScanSupport extends Logging {
isChangelogScan: Boolean): Boolean =
field.name == MetadataColumns.FILE_PATH.name() ||
field.name == MetadataColumns.SPEC_ID.name() ||
(!isChangelogScan && field.name == MetadataColumns.PARTITION_COLUMN_NAME) ||
(isChangelogScan && ChangelogMetadataColumnNames.contains(field.name))

private def deletesEmpty(deletes: java.util.List[_]): Boolean =
Expand Down Expand Up @@ -719,14 +723,23 @@ object IcebergScanSupport extends Logging {

private def toNativeScanTask(
task: FileScanTask,
partitionSchema: StructType): IcebergNativeScanTask = {
partitionSchema: StructType,
fieldIdsByName: Map[String, Int],
partitionType: Types.StructType): IcebergNativeScanTask = {
val file = task.file()
lazy val constants =
PartitionUtil.constantsMap(task, partitionType, SparkUtil.internalToSpark(_, _))
val values = partitionSchema.fields.map { field =>
CatalystTypeConverters.convertToScala(
constants.get(fieldIdsByName(field.name)),
field.dataType)
}
IcebergNativeScanTask(
file.location(),
task.start(),
task.length(),
file.fileSizeInBytes(),
metadataPartitionValues(file.location(), file.specId(), None, partitionSchema))
values.toSeq)
}

private def toNativeScanTask(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,8 +150,7 @@ case class NativeIcebergTableScanExec(
private def metadataPartitionValues(file: PartitionedFile): Seq[pb.ScalarValue] =
partitionSchema.fields.zipWithIndex.map { case (field, index) =>
NativeConverters
.convertExpr(
Literal.create(file.partitionValues.get(index, field.dataType), field.dataType))
.convertExpr(Literal(file.partitionValues.get(index, field.dataType), field.dataType))
.getLiteral
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -588,6 +588,115 @@ class AuronIcebergIntegrationSuite
}
}

Seq("parquet", "orc").foreach { format =>
test(s"iceberg native scan supports _partition on unpartitioned tables: $format") {
withTable("local.db.t_partition_metadata") {
sql(s"""CREATE TABLE local.db.t_partition_metadata (id INT) USING iceberg
|TBLPROPERTIES ('write.format.default' = '$format')""".stripMargin)
checkSparkAnswerAndOperator("SELECT id, _partition FROM local.db.t_partition_metadata")
sql("INSERT INTO local.db.t_partition_metadata VALUES (1), (2)")
val df =
checkSparkAnswerAndOperator("SELECT id, _partition FROM local.db.t_partition_metadata")
checkAnswer(df, Seq(Row(1, null), Row(2, null)))
checkSparkAnswerAndOperator("SELECT _partition FROM local.db.t_partition_metadata")
}
}

test(s"iceberg native scan supports _partition values and field projections: $format") {
withTable("local.db.t_partition_metadata") {
sql(s"""CREATE TABLE local.db.t_partition_metadata
|(id INT, p STRING, ts TIMESTAMP, d DATE, amount DECIMAL(8, 2)) USING iceberg
|PARTITIONED BY (p, bucket(4, id), days(ts), d, amount)
|TBLPROPERTIES ('write.format.default' = '$format')""".stripMargin)
sql("""INSERT INTO local.db.t_partition_metadata VALUES
|(1, 'east', TIMESTAMP '2026-01-01 12:00:00', DATE '2026-01-01', 12.34),
|(2, 'west', TIMESTAMP '2026-01-02 12:00:00', DATE '2026-01-02', 56.78),
|(3, NULL, NULL, NULL, NULL)""".stripMargin)

checkSparkAnswerAndOperator("SELECT _partition FROM local.db.t_partition_metadata")
checkSparkAnswerAndOperator("SELECT id, _partition FROM local.db.t_partition_metadata")
checkSparkAnswerAndOperator(
"SELECT _partition, _file, p, id, _spec_id FROM local.db.t_partition_metadata")
checkSparkAnswerAndOperator(
"SELECT _partition.p, _partition.id_bucket, _partition.ts_day, " +
"_partition.d, _partition.amount FROM local.db.t_partition_metadata")
val df = checkSparkAnswerAndOperator(
"SELECT id, _partition.p FROM local.db.t_partition_metadata")
checkAnswer(df, Seq(Row(1, "east"), Row(2, "west"), Row(3, null)))
}
}

test(s"iceberg native scan supports _partition across partition specs: $format") {
withTable("local.db.t_partition_metadata") {
sql(s"""CREATE TABLE local.db.t_partition_metadata
|(id INT, p STRING, ts TIMESTAMP) USING iceberg
|PARTITIONED BY (p, bucket(4, id))
|TBLPROPERTIES ('format-version' = '2', 'write.format.default' = '$format')
|""".stripMargin)
sql("""INSERT INTO local.db.t_partition_metadata VALUES
|(1, 'east', TIMESTAMP '2026-01-01 12:00:00'), (2, NULL, NULL)""".stripMargin)
sql("ALTER TABLE local.db.t_partition_metadata DROP PARTITION FIELD p")
sql("ALTER TABLE local.db.t_partition_metadata ADD PARTITION FIELD days(ts)")
sql("""INSERT INTO local.db.t_partition_metadata VALUES
|(3, 'west', TIMESTAMP '2026-01-03 12:00:00'), (4, NULL, NULL)""".stripMargin)

val df = checkSparkAnswerAndOperator(
"SELECT id, _partition, _spec_id FROM local.db.t_partition_metadata")
assert(df.collect().map(_.getInt(2)).distinct.length == 2)
checkSparkAnswerAndOperator("SELECT _partition FROM local.db.t_partition_metadata")
val fields = checkSparkAnswerAndOperator(
"SELECT id, _partition.p, _partition.ts_day FROM local.db.t_partition_metadata")
val rows = fields.collect().map(row => row.getInt(0) -> row).toMap
assert(rows(1).getString(1) == "east" && rows(1).isNullAt(2))
assert(rows(3).isNullAt(1) && !rows(3).isNullAt(2))
assert(rows(2).isNullAt(1) && rows(2).isNullAt(2))
assert(rows(4).isNullAt(1) && rows(4).isNullAt(2))
}
}
}

test("iceberg _partition scan falls back for unsupported partition types") {
withTable("local.db.t_partition_metadata") {
sql("""CREATE TABLE local.db.t_partition_metadata (id INT, p DECIMAL(20, 2))
|USING iceberg PARTITIONED BY (p)""".stripMargin)
sql("INSERT INTO local.db.t_partition_metadata VALUES (1, 12.34)")
val query = "SELECT _partition FROM local.db.t_partition_metadata"
var expected: Seq[Row] = Nil
withSQLConf("spark.auron.enabled" -> "false") {
expected = sql(query).collect().toSeq
}
withSQLConf("spark.auron.enabled" -> "true", "spark.auron.enable.iceberg.scan" -> "true") {
val df = sql(query)
checkAnswer(df, expected)
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
}
}
}

test("iceberg _partition scan falls back for row metadata and changelog tasks") {
withTable("local.db.t_partition_metadata") {
sql(
"CREATE TABLE local.db.t_partition_metadata (id INT, p STRING) " +
"USING iceberg PARTITIONED BY (p)")
sql("INSERT INTO local.db.t_partition_metadata VALUES (1, 'east')")
Seq(
"SELECT id, _partition, _pos FROM local.db.t_partition_metadata",
"SELECT id, _partition FROM local.db.t_partition_metadata.changes").foreach { query =>
var expected: Seq[Row] = Nil
withSQLConf("spark.auron.enabled" -> "false") {
expected = sql(query).collect().toSeq
}
withSQLConf(
"spark.auron.enabled" -> "true",
"spark.auron.enable.iceberg.scan" -> "true") {
val df = sql(query)
checkAnswer(df, expected)
assert(!df.queryExecution.executedPlan.toString().contains("NativeIcebergTableScan"))
}
}
}
}

test("iceberg native scan supports data columns with _file and _spec_id metadata columns") {
withTable("local.db.t4_metadata_mixed") {
sql("create table local.db.t4_metadata_mixed using iceberg as select 1 as id, 'a' as v")
Expand Down
Loading