From eba27282a997b8233e12bf6f1d795b44c051a240 Mon Sep 17 00:00:00 2001 From: zhangfengcdt Date: Wed, 7 Oct 2026 10:57:44 -0700 Subject: [PATCH 1/4] fix: fall back from native V2 scans when the plan reads input_file_name The native Iceberg and CSV V2 scans do not set InputFileBlockHolder, so input_file_name, input_file_block_start and input_file_block_length returned empty values. Fall back to Spark's reader as the V1 path does. Closes #6707 --- docs/source/user-guide/latest/iceberg.md | 1 + .../apache/comet/rules/CometScanRule.scala | 18 ++++++++-- .../comet/CometIcebergNativeSuite.scala | 36 +++++++++++++++++++ 3 files changed, 53 insertions(+), 2 deletions(-) diff --git a/docs/source/user-guide/latest/iceberg.md b/docs/source/user-guide/latest/iceberg.md index 4c3abe00216..fbdf81c3712 100644 --- a/docs/source/user-guide/latest/iceberg.md +++ b/docs/source/user-guide/latest/iceberg.md @@ -186,6 +186,7 @@ The following scenarios will fall back to the JVM Iceberg reader: - Encrypted tables with 192-bit data keys (no AES-192-GCM in the underlying crypto) - Delete files in a format other than Parquet or Puffin (Avro or ORC positional/equality deletes) - Tables backed by Avro or ORC data files (only Parquet is accelerated) +- Queries that use `input_file_name()`, `input_file_block_start()` or `input_file_block_length()` - Scans whose data or delete files span more than one S3 bucket (the native reader uses one object-store configuration per scan) - Tables partitioned on `BINARY` or `DECIMAL` (with precision >28) columns diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index b109fa158ed..3e8ddf936a9 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -111,7 +111,7 @@ case class CometScanRule(session: SparkSession) // is ever offered it. `transformV2Scan` applies the guard right after its contrib hook // declines. case scanExec: BatchScanExec => - transformV2Scan(scanExec) + transformV2Scan(fullPlan, scanExec) } plan.transform { @@ -379,7 +379,7 @@ case class CometScanRule(session: SparkSession) Some(CometScanExec(scanExec, session)) } - private def transformV2Scan(scanExec: BatchScanExec): SparkPlan = { + private def transformV2Scan(plan: SparkPlan, scanExec: BatchScanExec): SparkPlan = { // Give any optional, out-of-tree scan contrib (e.g. Lance) first crack at this V2 scan. On a // default build no contrib is registered, so this returns None and we proceed with Comet's @@ -398,6 +398,20 @@ case class CometScanRule(session: SparkSession) return withFallbackReason(scanExec, "Iceberg Metadata tables are not supported") } + // As in transformV1Scan: these expressions read InputFileBlockHolder, which the source's own + // reader sets per file. Comet's native V2 scans (Iceberg, CSV) do not, so they would return + // empty/default values (https://github.com/apache/datafusion-comet/issues/6707). + if (plan.exists(node => + node.expressions.exists(_.exists { + case _: InputFileName | _: InputFileBlockStart | _: InputFileBlockLength => true + case _ => false + }))) { + return withFallbackReason( + scanExec, + "Native V2 scan is not compatible with input_file_name, " + + "input_file_block_start, or input_file_block_length") + } + // NOTE: there is no blanket metadata-column guard here. Comet's built-in V2 paths handle // metadata columns individually -- the Iceberg path supports the ones in // `CometIcebergNativeScan.MetadataFieldIds` and rejects the rest, the CSV path rejects all -- diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala index f271f7d6ec2..8a8f1b53e55 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala @@ -166,6 +166,42 @@ class CometIcebergNativeSuite } } + // https://github.com/apache/datafusion-comet/issues/6707 + test("input_file_name and input_file_block_* fall back to Spark's Iceberg reader") { + assume(icebergAvailable, "Iceberg not available in classpath") + + withTempIcebergDir { warehouseDir => + withSQLConf( + "spark.sql.catalog.hadoop_catalog" -> "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.hadoop_catalog.type" -> "hadoop", + "spark.sql.catalog.hadoop_catalog.warehouse" -> warehouseDir.getAbsolutePath, + CometConf.COMET_ENABLED.key -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") { + + spark.sql("CREATE TABLE hadoop_catalog.db.input_file (id BIGINT) USING iceberg") + for (i <- 0 until 3) { + spark.sql( + "INSERT INTO hadoop_catalog.db.input_file " + + s"SELECT id FROM range(${i * 1000}, ${(i + 1) * 1000})") + } + val columns = "input_file_name(), input_file_block_start(), input_file_block_length(), id" + for (query <- Seq( + s"SELECT $columns FROM hadoop_catalog.db.input_file", + s"SELECT $columns FROM hadoop_catalog.db.input_file WHERE id >= 0")) { + val (_, cometPlan) = checkSparkAnswerAndFallbackReason(query, "input_file_name") + assert( + collectIcebergNativeScans(cometPlan).isEmpty, + s"Expected fallback to Spark but found a CometIcebergNativeScanExec. Plan:\n$cometPlan") + } + // Without these expressions the scan stays native + checkIcebergNativeScan("SELECT id FROM hadoop_catalog.db.input_file WHERE id >= 0") + + spark.sql("DROP TABLE hadoop_catalog.db.input_file") + } + } + } + test("filter pushdown - equality predicates") { assume(icebergAvailable, "Iceberg not available in classpath") From 474705a7df9a9cb1c140f219eec51f21966a16ab Mon Sep 17 00:00:00 2001 From: zhangfengcdt Date: Wed, 7 Oct 2026 18:02:10 -0700 Subject: [PATCH 2/4] refactor: reuse readsInputFileBlock in the V2 scan check #6703 added the shared helper for the same predicate. --- .../main/scala/org/apache/comet/rules/CometScanRule.scala | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index 5c02e12f47f..197f158c382 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -397,11 +397,7 @@ case class CometScanRule(session: SparkSession) // As in transformV1Scan: these expressions read InputFileBlockHolder, which the source's own // reader sets per file. Comet's native V2 scans (Iceberg, CSV) do not, so they would return // empty/default values (https://github.com/apache/datafusion-comet/issues/6707). - if (plan.exists(node => - node.expressions.exists(_.exists { - case _: InputFileName | _: InputFileBlockStart | _: InputFileBlockLength => true - case _ => false - }))) { + if (CometScanRule.readsInputFileBlock(plan)) { return withFallbackReason( scanExec, "Native V2 scan is not compatible with input_file_name, " + From 5502bef0c7de352be9ce1139596e7ffbbbfa675a Mon Sep 17 00:00:00 2001 From: zhangfengcdt Date: Thu, 8 Oct 2026 10:52:00 -0700 Subject: [PATCH 3/4] test: cover the native CSV V2 scan's input_file_name fallback Document the fallback in the CSV section of datasources.md. --- docs/source/user-guide/latest/datasources.md | 4 ++- .../comet/csv/CometCsvNativeReadSuite.scala | 27 +++++++++++++++++++ 2 files changed, 30 insertions(+), 1 deletion(-) diff --git a/docs/source/user-guide/latest/datasources.md b/docs/source/user-guide/latest/datasources.md index f351aff5b6f..4fdfb11ceaf 100644 --- a/docs/source/user-guide/latest/datasources.md +++ b/docs/source/user-guide/latest/datasources.md @@ -39,7 +39,9 @@ See the [Iceberg Guide] and [Iceberg Writes](iceberg-writes.md) for more informa Comet provides experimental Rust-based CSV scan support. When `spark.comet.scan.csv.v2.enabled` is enabled, CSV files are read in Rust for improved performance. This feature is experimental and performance benefits are workload-dependent. Only Spark's DataSource V2 CSV scan is accelerated, and Spark reads CSV through the V1 API by -default, so also remove `csv` from `spark.sql.sources.useV1SourceList`. +default, so also remove `csv` from `spark.sql.sources.useV1SourceList`. Queries that use +`input_file_name()`, `input_file_block_start()` or `input_file_block_length()` fall back to Spark's CSV +reader, which sets the values these functions read. Alternatively, when `spark.comet.convert.csv.enabled` is enabled, data from Spark's CSV reader is immediately converted into Arrow format, allowing the Comet pipeline to take over after that. diff --git a/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala b/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala index 6e743cfc90d..7ab03c4ee09 100644 --- a/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala @@ -147,4 +147,31 @@ class CometCsvNativeReadSuite extends CometTestBase { } } } + + // https://github.com/apache/datafusion-comet/issues/6707. These functions read + // InputFileBlockHolder, which the native CSV scan does not set. + test("Native csv read - input_file_name falls back to Spark's reader") { + withTempPath { dir => + spark.range(30).repartition(3).write.csv(dir.toString) + withSQLConf( + CometConf.COMET_CSV_V2_NATIVE_ENABLED.key -> "true", + SQLConf.USE_V1_SOURCE_LIST.key -> "") { + val read = spark.read.schema("id LONG").csv(dir.toString) + val df = read.selectExpr( + "input_file_name()", + "input_file_block_start()", + "input_file_block_length()", + "id") + val (_, plan) = checkSparkAnswerAndFallbackReason(df, "input_file_name") + assert( + collect(plan) { case s: CometCsvNativeScanExec => s }.isEmpty, + s"Expected Spark's CSV reader but found CometCsvNativeScanExec. Plan:\n$plan") + // Without these functions the scan stays native + val (_, nativePlan) = checkSparkAnswer(read.where("id >= 0")) + assert( + collect(nativePlan) { case s: CometCsvNativeScanExec => s }.nonEmpty, + s"Expected CometCsvNativeScanExec. Plan:\n$nativePlan") + } + } + } } From e820753d6af6b77a5a0da1ba4500d193d435b154 Mon Sep 17 00:00:00 2001 From: zhangfengcdt Date: Fri, 9 Oct 2026 08:45:16 -0700 Subject: [PATCH 4/4] test: move the Iceberg input_file_name test to a SQL fixture Match the full fallback reason there and in the CSV test. --- .../sql-tests/iceberg/input_file_name.sql | 61 +++++++++++++++++++ .../comet/CometIcebergNativeSuite.scala | 36 ----------- .../comet/csv/CometCsvNativeReadSuite.scala | 4 +- 3 files changed, 64 insertions(+), 37 deletions(-) create mode 100644 spark/src/test/resources/sql-tests/iceberg/input_file_name.sql diff --git a/spark/src/test/resources/sql-tests/iceberg/input_file_name.sql b/spark/src/test/resources/sql-tests/iceberg/input_file_name.sql new file mode 100644 index 00000000000..a473e720ae9 --- /dev/null +++ b/spark/src/test/resources/sql-tests/iceberg/input_file_name.sql @@ -0,0 +1,61 @@ +-- Licensed to the Apache Software Foundation (ASF) under one +-- or more contributor license agreements. See the NOTICE file +-- distributed with this work for additional information +-- regarding copyright ownership. The ASF licenses this file +-- to you under the Apache License, Version 2.0 (the +-- "License"); you may not use this file except in compliance +-- with the License. You may obtain a copy of the License at +-- +-- http://www.apache.org/licenses/LICENSE-2.0 +-- +-- Unless required by applicable law or agreed to in writing, +-- software distributed under the License is distributed on an +-- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +-- KIND, either express or implied. See the License for the +-- specific language governing permissions and limitations +-- under the License. + +-- input_file_name, input_file_block_start and input_file_block_length read InputFileBlockHolder, +-- which Iceberg's Spark reader sets as it opens each data file. The native Iceberg scan does not, +-- so a plan that uses them falls back to Spark's reader +-- (https://github.com/apache/datafusion-comet/issues/6707). + +-- Native Iceberg scan is unsupported on Spark 4.2: no Iceberg spark-runtime is published for +-- 4.2 yet. See https://github.com/apache/datafusion-comet/issues/4969. +-- MaxSparkVersion: 4.1 + +-- Config: spark.sql.catalog.test_cat=org.apache.iceberg.spark.SparkCatalog +-- Config: spark.sql.catalog.test_cat.type=hadoop +-- Config: spark.sql.catalog.test_cat.warehouse=/tmp/comet-iceberg-sql-test +-- Config: spark.comet.enabled=true +-- Config: spark.comet.exec.enabled=true +-- Config: spark.comet.scan.icebergNative.enabled=true + +statement +DROP TABLE IF EXISTS test_cat.db.input_file + +statement +CREATE TABLE test_cat.db.input_file (id BIGINT) USING iceberg + +-- Three inserts write three data files +statement +INSERT INTO test_cat.db.input_file SELECT id FROM range(0, 1000) + +statement +INSERT INTO test_cat.db.input_file SELECT id FROM range(1000, 2000) + +statement +INSERT INTO test_cat.db.input_file SELECT id FROM range(2000, 3000) + +query expect_fallback(Native V2 scan is not compatible with input_file_name) +SELECT input_file_name(), input_file_block_start(), input_file_block_length(), id +FROM test_cat.db.input_file + +-- A Comet filter above the scan +query expect_fallback(Native V2 scan is not compatible with input_file_name) +SELECT input_file_name(), input_file_block_start(), input_file_block_length(), id +FROM test_cat.db.input_file WHERE id >= 0 + +-- Without these functions the scan stays native +query +SELECT id FROM test_cat.db.input_file WHERE id >= 0 diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala index 8a8f1b53e55..f271f7d6ec2 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala @@ -166,42 +166,6 @@ class CometIcebergNativeSuite } } - // https://github.com/apache/datafusion-comet/issues/6707 - test("input_file_name and input_file_block_* fall back to Spark's Iceberg reader") { - assume(icebergAvailable, "Iceberg not available in classpath") - - withTempIcebergDir { warehouseDir => - withSQLConf( - "spark.sql.catalog.hadoop_catalog" -> "org.apache.iceberg.spark.SparkCatalog", - "spark.sql.catalog.hadoop_catalog.type" -> "hadoop", - "spark.sql.catalog.hadoop_catalog.warehouse" -> warehouseDir.getAbsolutePath, - CometConf.COMET_ENABLED.key -> "true", - CometConf.COMET_EXEC_ENABLED.key -> "true", - CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") { - - spark.sql("CREATE TABLE hadoop_catalog.db.input_file (id BIGINT) USING iceberg") - for (i <- 0 until 3) { - spark.sql( - "INSERT INTO hadoop_catalog.db.input_file " + - s"SELECT id FROM range(${i * 1000}, ${(i + 1) * 1000})") - } - val columns = "input_file_name(), input_file_block_start(), input_file_block_length(), id" - for (query <- Seq( - s"SELECT $columns FROM hadoop_catalog.db.input_file", - s"SELECT $columns FROM hadoop_catalog.db.input_file WHERE id >= 0")) { - val (_, cometPlan) = checkSparkAnswerAndFallbackReason(query, "input_file_name") - assert( - collectIcebergNativeScans(cometPlan).isEmpty, - s"Expected fallback to Spark but found a CometIcebergNativeScanExec. Plan:\n$cometPlan") - } - // Without these expressions the scan stays native - checkIcebergNativeScan("SELECT id FROM hadoop_catalog.db.input_file WHERE id >= 0") - - spark.sql("DROP TABLE hadoop_catalog.db.input_file") - } - } - } - test("filter pushdown - equality predicates") { assume(icebergAvailable, "Iceberg not available in classpath") diff --git a/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala b/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala index 7ab03c4ee09..eeb809764be 100644 --- a/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala +++ b/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala @@ -162,7 +162,9 @@ class CometCsvNativeReadSuite extends CometTestBase { "input_file_block_start()", "input_file_block_length()", "id") - val (_, plan) = checkSparkAnswerAndFallbackReason(df, "input_file_name") + val (_, plan) = checkSparkAnswerAndFallbackReason( + df, + "Native V2 scan is not compatible with input_file_name") assert( collect(plan) { case s: CometCsvNativeScanExec => s }.isEmpty, s"Expected Spark's CSV reader but found CometCsvNativeScanExec. Plan:\n$plan")