Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
4 changes: 3 additions & 1 deletion docs/source/user-guide/latest/datasources.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions docs/source/user-guide/latest/iceberg.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
14 changes: 12 additions & 2 deletions spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand All @@ -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)) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This also changes the native CSV V2 scan, since it goes through the same check. I tried it locally. With spark.comet.scan.csv.v2.enabled=true and an empty spark.sql.sources.useV1SourceList, a CSV read that selects input_file_name() returns different results from Spark without this check. With it, the scan falls back and matches. Could we add a case like that to CometCsvNativeReadSuite, with a few small CSV files and a check for the new fallback reason? The CSV scan is off by default, so if a later change moved this check into the Iceberg branch, nothing would catch it. Could we also add a sentence to the CSV section of datasources.md saying the native CSV scan falls back for these functions, the way iceberg.md now does?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure! I have add both. CometCsvNativeReadSuite now has "input_file_name falls back to Spark's reader". It fails with mismatched results if the V2 check is removed.

The CSV section of datasources.md now notes the fallback.

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 --
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
}
}
Loading