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/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 c7c9a31909d..197f158c382 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 { @@ -375,7 +375,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 @@ -394,6 +394,16 @@ 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 (CometScanRule.readsInputFileBlock(plan)) { + 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/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/csv/CometCsvNativeReadSuite.scala b/spark/src/test/scala/org/apache/comet/csv/CometCsvNativeReadSuite.scala index 6e743cfc90d..eeb809764be 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,33 @@ 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, + "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") + // 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") + } + } + } }