andygrove commented on code in PR #6784:
URL: https://github.com/apache/datafusion-comet/pull/6784#discussion_r4223872700


##########
spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala:
##########
@@ -458,6 +459,147 @@ case class CometScanRule(session: SparkSession)
           withFallbackReasons(scanExec, fallbackReasons.toSet)
         }
 
+      case scan: TextScan if COMET_TEXT_V2_NATIVE_ENABLED.get() =>
+        if (scanExec.output.exists(_.isMetadataCol)) {
+          return withFallbackReason(
+            scanExec,
+            "Metadata columns are not supported for Text V2 scans")
+        }
+        // A native scan is useless without native execution: with exec 
disabled, CometExecRule
+        // leaves this CometBatchScanExec unconverted and it fails at runtime. 
Require native exec.
+        if (!COMET_EXEC_ENABLED.get()) {
+          return withFallbackReason(
+            scanExec,
+            s"Native Text scan requires ${COMET_EXEC_ENABLED.key} to be 
enabled")
+        }
+        // input_file_name/_block_start/_block_length read a thread-local that 
Spark's FileScanRDD
+        // sets; the native DataFusion scan does not, so they would return "" 
/ -1. Same guard as
+        // the native Parquet path.
+        if (plan.exists(node =>

Review Comment:
   This inline check is the same as `CometScanRule.readsInputFileBlock(plan)`, 
which the Parquet gate calls. Could the Text branch call the helper so the two 
cannot drift?



##########
spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala:
##########
@@ -458,6 +459,147 @@ case class CometScanRule(session: SparkSession)
           withFallbackReasons(scanExec, fallbackReasons.toSet)
         }
 
+      case scan: TextScan if COMET_TEXT_V2_NATIVE_ENABLED.get() =>
+        if (scanExec.output.exists(_.isMetadataCol)) {
+          return withFallbackReason(
+            scanExec,
+            "Metadata columns are not supported for Text V2 scans")
+        }
+        // A native scan is useless without native execution: with exec 
disabled, CometExecRule
+        // leaves this CometBatchScanExec unconverted and it fails at runtime. 
Require native exec.
+        if (!COMET_EXEC_ENABLED.get()) {
+          return withFallbackReason(
+            scanExec,
+            s"Native Text scan requires ${COMET_EXEC_ENABLED.key} to be 
enabled")
+        }
+        // input_file_name/_block_start/_block_length read a thread-local that 
Spark's FileScanRDD
+        // sets; the native DataFusion scan does not, so they would return "" 
/ -1. Same guard as
+        // the native Parquet path.
+        if (plan.exists(node =>
+            node.expressions.exists(_.exists {
+              case _: InputFileName | _: InputFileBlockStart | _: 
InputFileBlockLength => true
+              case _ => false
+            }))) {
+          return withFallbackReason(
+            scanExec,
+            "Native Text scan is not compatible with input_file_name, " +
+              "input_file_block_start, or input_file_block_length")
+        }
+        val fallbackReasons = new ListBuffer[String]()
+        val schemaSupported =
+          CometBatchScanExec.isSchemaSupported(scan.readDataSchema, 
fallbackReasons)
+        if (!schemaSupported) {
+          fallbackReasons += s"Schema ${scan.readDataSchema} is not supported"
+        }
+        // Spark's text source allows only a single data column 
(TextScan.verifyReadSchema). A
+        // user-forced multi-column read schema is invalid; fall back so Spark 
raises its structured
+        // AnalysisException rather than letting the native reader fail with 
an opaque error.
+        if (scan.readDataSchema.size > 1) {
+          fallbackReasons += "Text data source supports only a single column"
+        }
+        // Partition columns are not implemented for native Text scans (unlike 
the native CSV path,
+        // which forwards them). Fall back to Spark for partitioned reads 
rather than drop them.
+        if (scan.readPartitionSchema.nonEmpty) {
+          fallbackReasons += "Comet does not support partition columns in 
native Text scans"
+        }
+        // The native reader loads each file whole. Spark splits a large 
uncompressed text file
+        // into byte-range partitions (TextScan.isSplitable), and each split 
would then re-read and
+        // buffer the whole file with no memory-pool reservation. Restrict 
native execution to
+        // unsplit files (the small build-side lookups this targets); Spark 
handles large files.
+        val hasSplitFile = scanExec.inputPartitions.exists {
+          case fp: FilePartition => fp.files.exists(_.start != 0)
+          case _ => false
+        }
+        if (hasSplitFile) {
+          fallbackReasons += "Comet native Text scan does not support files 
split into byte " +
+            "ranges (large files)"
+        }
+        // Comet's native plan does not carry a compression codec, so the 
native reader would feed
+        // raw compressed bytes to the line splitter. Spark auto-detects a 
codec by file extension
+        // (e.g. .gz, .bz2). Fall back when any input file is compressed. Use 
the options-aware
+        // Hadoop conf so a codec registered via a per-read option is honored 
(as Spark does).
+        val hadoopConf = 
session.sessionState.newHadoopConfWithOptions(scan.options.asScala.toMap)
+        val codecFactory = new 
org.apache.hadoop.io.compress.CompressionCodecFactory(hadoopConf)
+        val hasCompressedFile = scan.fileIndex.inputFiles.exists { path =>
+          codecFactory.getCodec(new org.apache.hadoop.fs.Path(path)) != null
+        }
+        if (hasCompressedFile) {
+          fallbackReasons += "Comet does not support compressed text files"

Review Comment:
   I think the compressed-file check misses two extensions on Spark 4.1 and 
later. `HadoopLineRecordReader` and `WholeTextFileRecordReader` pick the codec 
through `HadoopCodecStreams`, which also treats `.gzip` and `.zstd` as 
compressed. `CompressionCodecFactory.getCodec` returns null for both, so these 
files stay native and the rows are the raw compressed bytes.
   
   I reproduced this on Spark 4.1.3. A gzipped `l1\nl2\nl3\n` saved as 
`b.txt.gzip` returns `l1, l2, l3` from Spark, and the native scan returns the 
gzip bytes with no error. A `.zstd` file does the same, and so does 
`wholetext`. On Spark 3.5 Spark reads these files raw as well, so the two sides 
agree there and the check has to be version aware.
   
   Would a small shim next to `ShimFileFormat` work? It could call 
`HadoopCodecStreams.getDecompressionCodec` in the `spark-4.1+` source set and 
`CompressionCodecFactory` everywhere else. Tests for a `.gzip` file and a 
`.zstd` file would lock it in.



##########
spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala:
##########
@@ -458,6 +459,147 @@ case class CometScanRule(session: SparkSession)
           withFallbackReasons(scanExec, fallbackReasons.toSet)
         }
 
+      case scan: TextScan if COMET_TEXT_V2_NATIVE_ENABLED.get() =>
+        if (scanExec.output.exists(_.isMetadataCol)) {
+          return withFallbackReason(
+            scanExec,
+            "Metadata columns are not supported for Text V2 scans")
+        }
+        // A native scan is useless without native execution: with exec 
disabled, CometExecRule
+        // leaves this CometBatchScanExec unconverted and it fails at runtime. 
Require native exec.
+        if (!COMET_EXEC_ENABLED.get()) {
+          return withFallbackReason(
+            scanExec,
+            s"Native Text scan requires ${COMET_EXEC_ENABLED.key} to be 
enabled")
+        }
+        // input_file_name/_block_start/_block_length read a thread-local that 
Spark's FileScanRDD
+        // sets; the native DataFusion scan does not, so they would return "" 
/ -1. Same guard as
+        // the native Parquet path.
+        if (plan.exists(node =>
+            node.expressions.exists(_.exists {
+              case _: InputFileName | _: InputFileBlockStart | _: 
InputFileBlockLength => true
+              case _ => false
+            }))) {
+          return withFallbackReason(
+            scanExec,
+            "Native Text scan is not compatible with input_file_name, " +
+              "input_file_block_start, or input_file_block_length")
+        }
+        val fallbackReasons = new ListBuffer[String]()
+        val schemaSupported =
+          CometBatchScanExec.isSchemaSupported(scan.readDataSchema, 
fallbackReasons)
+        if (!schemaSupported) {
+          fallbackReasons += s"Schema ${scan.readDataSchema} is not supported"
+        }
+        // Spark's text source allows only a single data column 
(TextScan.verifyReadSchema). A
+        // user-forced multi-column read schema is invalid; fall back so Spark 
raises its structured
+        // AnalysisException rather than letting the native reader fail with 
an opaque error.
+        if (scan.readDataSchema.size > 1) {
+          fallbackReasons += "Text data source supports only a single column"
+        }
+        // Partition columns are not implemented for native Text scans (unlike 
the native CSV path,
+        // which forwards them). Fall back to Spark for partitioned reads 
rather than drop them.
+        if (scan.readPartitionSchema.nonEmpty) {
+          fallbackReasons += "Comet does not support partition columns in 
native Text scans"
+        }
+        // The native reader loads each file whole. Spark splits a large 
uncompressed text file
+        // into byte-range partitions (TextScan.isSplitable), and each split 
would then re-read and
+        // buffer the whole file with no memory-pool reservation. Restrict 
native execution to
+        // unsplit files (the small build-side lookups this targets); Spark 
handles large files.
+        val hasSplitFile = scanExec.inputPartitions.exists {
+          case fp: FilePartition => fp.files.exists(_.start != 0)
+          case _ => false
+        }
+        if (hasSplitFile) {
+          fallbackReasons += "Comet native Text scan does not support files 
split into byte " +
+            "ranges (large files)"
+        }
+        // Comet's native plan does not carry a compression codec, so the 
native reader would feed
+        // raw compressed bytes to the line splitter. Spark auto-detects a 
codec by file extension
+        // (e.g. .gz, .bz2). Fall back when any input file is compressed. Use 
the options-aware
+        // Hadoop conf so a codec registered via a per-read option is honored 
(as Spark does).
+        val hadoopConf = 
session.sessionState.newHadoopConfWithOptions(scan.options.asScala.toMap)
+        val codecFactory = new 
org.apache.hadoop.io.compress.CompressionCodecFactory(hadoopConf)
+        val hasCompressedFile = scan.fileIndex.inputFiles.exists { path =>
+          codecFactory.getCodec(new org.apache.hadoop.fs.Path(path)) != null
+        }
+        if (hasCompressedFile) {
+          fallbackReasons += "Comet does not support compressed text files"
+        }
+        // Spark skips unreadable/missing files when these are set; the native 
reader's DataFusion
+        // FileStream fails the whole scan instead (OnError::Fail). Fall back 
to match Spark.
+        if (SQLConf.get.ignoreCorruptFiles || scan.options.getBoolean(

Review Comment:
   The native reader also ignores 
`mapreduce.input.linerecordreader.line.maxlength`. Spark's line reader skips 
every line at or over that length. With the limit set to 5 and a file of `ab`, 
`abcdefghij` and `xy`, Spark returns two rows and the native scan returns 
three. I saw the same on Spark 3.5 and 4.1. The key can come from 
`spark.hadoop.*` or `core-site.xml` as well as a read option, so it is not only 
a per-read setting.
   
   Should the rule fall back when that key is set in `hadoopConf`, next to the 
`ignoreCorruptFiles` and `ignoreMissingFiles` checks? A test with the option 
set would keep it covered.



##########
spark/src/test/scala/org/apache/comet/text/CometTextNativeReadSuite.scala:
##########
@@ -0,0 +1,279 @@
+/*
+ * 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.
+ */
+
+package org.apache.comet.text
+
+import java.io.File
+import java.nio.charset.StandardCharsets
+import java.nio.file.Files
+
+import org.apache.spark.sql.CometTestBase
+import org.apache.spark.sql.functions.input_file_name
+import org.apache.spark.sql.internal.SQLConf
+
+import org.apache.comet.CometConf
+
+class CometTextNativeReadSuite extends CometTestBase {
+  private val TEST_TEXT_PATH = "src/test/resources/test-data/text-test-1.txt"
+
+  private def withTextConf(f: => Unit): Unit = {
+    withSQLConf(
+      CometConf.COMET_TEXT_V2_NATIVE_ENABLED.key -> "true",
+      SQLConf.USE_V1_SOURCE_LIST.key -> "")(f)
+  }
+
+  test("native text read - lines to value column") {
+    withTextConf {
+      val df = spark.read.text(TEST_TEXT_PATH)
+      checkSparkAnswerAndOperator(df)
+    }
+  }
+
+  test("native text read - count") {
+    withTextConf {
+      val df = spark.read.text(TEST_TEXT_PATH).selectExpr("count(*) as c")
+      checkSparkAnswerAndOperator(df)
+    }
+  }
+
+  test("native text read - filter and projection") {
+    withTextConf {
+      val df = spark.read.text(TEST_TEXT_PATH).where("value like '%bar%'")
+      checkSparkAnswerAndOperator(df)
+    }
+  }
+
+  test("native text read - preserves line order within a file") {
+    withTempDir { dir =>
+      val file = new File(dir, "ordered.txt")
+      val lines = (1 to 500).map(i => s"row-$i")
+      Files.write(file.toPath, 
lines.mkString("\n").getBytes(StandardCharsets.UTF_8))
+      withTextConf {
+        val got = 
spark.read.text(file.getAbsolutePath).collect().map(_.getString(0)).toSeq

Review Comment:
   This test only calls `collect()`, so it would still pass if the scan fell 
back to Spark. Could it also assert that the executed plan contains 
`CometTextNativeScanExec`?



##########
docs/source/user-guide/latest/datasources.md:
##########
@@ -49,6 +49,17 @@ converted into Arrow format, allowing the Comet pipeline to 
take over after that
 Comet does not provide a Rust-based JSON scan, but when 
`spark.comet.convert.json.enabled` is enabled, data is immediately
 converted into Arrow format, allowing the Comet pipeline to take over after 
that.
 
+### Text
+
+Comet provides experimental Rust-based text scan support (a single `value: 
string` column). When
+`spark.comet.scan.text.v2.enabled` is enabled, text files are read in Rust. 
This feature is experimental and
+performance benefits are workload-dependent. Only Spark's DataSource V2 text 
scan is accelerated, and Spark reads
+text through the V1 API by default, so also remove `text` from 
`spark.sql.sources.useV1SourceList`. The native scan
+falls back to Spark for cases it does not handle -- including partitioned 
tables, compressed files, files large
+enough that Spark splits them into byte ranges, `wholetext` files larger than 
roughly 2GB (capped at
+`min(spark.sql.files.maxPartitionBytes, 2GB)`), unsupported filesystem 
schemes, and `input_file_name()` or

Review Comment:
   Two details here don't match the code. The `wholetext` cap in 
`CometScanRule` is `min(spark.sql.files.maxPartitionBytes, Int.MaxValue / 3)`, 
which is 128 MB by default, so "roughly 2GB" is misleading. And only reads of 
partition columns fall back. Selecting just `value` from a partitioned 
directory stays native.
   
   Text files are also the most likely way for ill-formed UTF-8 to reach Comet. 
The reader replaces it with U+FFFD, so byte-level functions and equality can 
differ from Spark on those rows. For the bytes `61 FF 62`, `value = 'a\uFFFDb'` 
is true in Comet and false in Spark. Could this section link to "Strings with 
non-UTF-8 bytes" in the compatibility guide?



##########
native/core/src/execution/operators/text_scan.rs:
##########
@@ -0,0 +1,614 @@
+// 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.
+
+//! Native Text file scan, mirroring Spark's `text` data source.
+//!
+//! Spark's text reader produces a single `value: string` column. With the 
default options each
+//! line of the file is one row (split on `\n`, `\r\n`, or `\r`, with the 
terminator stripped).
+//! With `wholetext=true` the whole file is a single row. A custom `lineSep` 
splits on exactly that
+//! byte sequence instead of the universal-newline set. See the 
`TextOptions`/`HadoopFileLinesReader`
+//! logic in Spark's `org.apache.spark.sql.execution.datasources.text` package.
+
+use crate::execution::operators::ExecutionError;
+use arrow::array::StringBuilder;
+use arrow::datatypes::SchemaRef;
+use arrow::record_batch::{RecordBatch, RecordBatchOptions};
+use datafusion::common::Result;
+use datafusion::datasource::object_store::ObjectStoreUrl;
+use datafusion::physical_expr::projection::ProjectionExprs;
+use datafusion::physical_plan::metrics::{BaselineMetrics, 
ExecutionPlanMetricsSet};
+use datafusion::physical_plan::DisplayFormatType;
+use datafusion_comet_common::decode_utf8_spark_lossy;
+use datafusion_comet_proto::spark_operator::TextOptions;
+use datafusion_datasource::file::FileSource;
+use datafusion_datasource::file_compression_type::FileCompressionType;
+use datafusion_datasource::file_groups::FileGroup;
+use datafusion_datasource::file_scan_config::{FileScanConfig, 
FileScanConfigBuilder};
+use datafusion_datasource::file_stream::{FileOpenFuture, FileOpener};
+use datafusion_datasource::projection::{ProjectionOpener, SplitProjection};
+use datafusion_datasource::source::DataSourceExec;
+use datafusion_datasource::{as_file_source, PartitionedFile, TableSchema};
+use futures::StreamExt;
+use object_store::{ObjectStore, ObjectStoreExt};
+use std::borrow::Cow;
+use std::fmt;
+use std::io::Read;
+use std::sync::Arc;
+
+/// Batch size used when the execution framework does not set one before 
opening the file.
+const DEFAULT_TEXT_BATCH_SIZE: usize = 8192;
+
+/// UTF-8 byte-order mark. Spark's line reader strips a leading one; wholetext 
keeps it.
+const UTF8_BOM: &[u8] = &[0xEF, 0xBB, 0xBF];
+
+/// Cap on a batch's decoded `value` buffer bytes, kept well below Arrow's 2GB 
32-bit offset limit
+/// so a batch of many rows (or a few long lines) cannot overflow. Measured 
against the decoded
+/// string length, since a lossy UTF-8 decode can expand ill-formed bytes up 
to 3x.
+const MAX_BATCH_VALUE_BYTES: usize = 1 << 30;
+
+pub fn init_text_datasource_exec(
+    object_store_url: ObjectStoreUrl,
+    file_groups: Vec<Vec<PartitionedFile>>,
+    data_schema: SchemaRef,
+    _partition_schema: SchemaRef,
+    projection_vector: Vec<usize>,
+    text_options: &TextOptions,
+) -> Result<Arc<DataSourceExec>, ExecutionError> {
+    let text_source = TextSource::new(data_schema, text_options);
+
+    let file_groups = file_groups.into_iter().map(FileGroup::new).collect();
+
+    let file_scan_config = FileScanConfigBuilder::new(object_store_url, 
Arc::new(text_source))
+        .with_file_groups(file_groups)
+        .with_projection_indices(Some(projection_vector))?
+        .build();
+
+    Ok(DataSourceExec::from_data_source(file_scan_config))
+}
+
+/// A [`FileSource`] for Spark's text format. The file schema is always the 
single `value` column;
+/// projection (and partition columns, if any) is handled by 
[`ProjectionOpener`], exactly as the
+/// CSV source does.
+#[derive(Debug, Clone)]
+struct TextSource {
+    table_schema: TableSchema,
+    projection: SplitProjection,
+    batch_size: Option<usize>,
+    metrics: ExecutionPlanMetricsSet,
+    whole_text: bool,
+    // The raw separator bytes when a custom `lineSep` is set; None means 
universal newline
+    // (`\n`, `\r\n`, `\r`).
+    line_sep: Option<Vec<u8>>,
+}
+
+impl TextSource {
+    fn new(file_schema: SchemaRef, options: &TextOptions) -> Self {
+        let table_schema: TableSchema = file_schema.into();
+        // `line_sep` arrives as the exact separator bytes Spark computed 
(`lineSeparatorInRead`,
+        // already encoded with the file's charset), so we split on them 
directly.
+        let line_sep = options.line_sep.clone().filter(|sep| !sep.is_empty());
+        Self {
+            projection: SplitProjection::unprojected(&table_schema),
+            table_schema,
+            batch_size: None,
+            metrics: ExecutionPlanMetricsSet::new(),
+            whole_text: options.whole_text,
+            line_sep,
+        }
+    }
+}
+
+impl From<TextSource> for Arc<dyn FileSource> {
+    fn from(source: TextSource) -> Self {
+        as_file_source(source)
+    }
+}
+
+impl FileSource for TextSource {
+    fn create_file_opener(
+        &self,
+        object_store: Arc<dyn ObjectStore>,
+        base_config: &FileScanConfig,
+        partition_index: usize,
+    ) -> Result<Arc<dyn FileOpener>> {
+        // The inner opener must emit batches whose schema is the file schema 
projected by
+        // `file_indices` (that is what ProjectionOpener feeds its projector). 
For text that is
+        // either the single `value` column or, for a projection-less count, 
an empty schema.
+        let file_schema = self.table_schema.file_schema();
+        let projected_file_schema = 
Arc::new(file_schema.project(&self.projection.file_indices)?);
+        let opener = Arc::new(TextOpener {
+            projected_file_schema,
+            file_compression_type: base_config.file_compression_type,
+            object_store,
+            partition_index,
+            batch_size: self.batch_size.unwrap_or(DEFAULT_TEXT_BATCH_SIZE),
+            metrics: self.metrics.clone(),
+            whole_text: self.whole_text,
+            line_sep: self.line_sep.clone(),
+        }) as Arc<dyn FileOpener>;
+        ProjectionOpener::try_new(self.projection.clone(), opener, file_schema)
+    }
+
+    fn table_schema(&self) -> &TableSchema {
+        &self.table_schema
+    }
+
+    fn with_batch_size(&self, batch_size: usize) -> Arc<dyn FileSource> {
+        let mut conf = self.clone();
+        conf.batch_size = Some(batch_size);
+        Arc::new(conf)
+    }
+
+    fn try_pushdown_projection(
+        &self,
+        projection: &ProjectionExprs,
+    ) -> Result<Option<Arc<dyn FileSource>>> {
+        let mut source = self.clone();
+        let new_projection = self.projection.source.try_merge(projection)?;
+        source.projection = 
SplitProjection::new(self.table_schema.file_schema(), &new_projection);
+        Ok(Some(Arc::new(source)))
+    }
+
+    fn projection(&self) -> Option<&ProjectionExprs> {
+        Some(&self.projection.source)
+    }
+
+    fn metrics(&self) -> &ExecutionPlanMetricsSet {
+        &self.metrics
+    }
+
+    fn file_type(&self) -> &str {
+        "text"
+    }
+
+    // A custom line separator can straddle a byte-range boundary, and we read 
whole files rather
+    // than boundary-aligned ranges, so this source must not let DataFusion 
repartition files by
+    // byte range.
+    fn supports_repartitioning(&self) -> bool {
+        false
+    }
+
+    fn fmt_extra(&self, t: DisplayFormatType, f: &mut fmt::Formatter) -> 
fmt::Result {
+        match t {
+            DisplayFormatType::Default | DisplayFormatType::Verbose => {
+                write!(f, ", whole_text={}", self.whole_text)
+            }
+            DisplayFormatType::TreeRender => Ok(()),
+        }
+    }
+
+    fn apply_expressions(
+        &self,
+        f: &mut dyn FnMut(
+            &Arc<dyn datafusion::physical_plan::PhysicalExpr>,
+        ) -> Result<datafusion::common::tree_node::TreeNodeRecursion>,
+    ) -> Result<datafusion::common::tree_node::TreeNodeRecursion> {
+        
datafusion::physical_plan::apply_expression_roots(self.projection.source.iter(),
 f)
+    }
+}
+
+/// A [`FileOpener`] that reads a whole text file and yields the projected 
file schema in batches.
+struct TextOpener {
+    // The file schema projected by `file_indices`: the single `value` column, 
or an empty schema
+    // for a projection-less count.
+    projected_file_schema: SchemaRef,
+    file_compression_type: FileCompressionType,
+    object_store: Arc<dyn ObjectStore>,
+    partition_index: usize,
+    batch_size: usize,
+    metrics: ExecutionPlanMetricsSet,
+    whole_text: bool,
+    line_sep: Option<Vec<u8>>,
+}
+
+impl FileOpener for TextOpener {
+    fn open(&self, partitioned_file: PartitionedFile) -> 
Result<FileOpenFuture> {
+        let store = Arc::clone(&self.object_store);
+        let projected_file_schema = Arc::clone(&self.projected_file_schema);
+        let file_compression_type = self.file_compression_type;
+        let batch_size = self.batch_size;
+        let whole_text = self.whole_text;
+        let line_sep = self.line_sep.clone();
+        let baseline_metrics = BaselineMetrics::new(&self.metrics, 
self.partition_index);
+
+        Ok(Box::pin(async move {
+            let range = partitioned_file.range.clone();
+            let location = partitioned_file.object_meta.location;
+            let raw = store.get(&location).await?.bytes().await?;

Review Comment:
   This reads the whole object into memory. After that `split_lines` builds an 
entry for every line, and `build_batches` builds every batch before the stream 
returns its first one. None of it is reserved with the memory pool.
   
   I measured a 119 MB file with 5M short lines on Spark 4.1, running 
`max(length(value))` over one unsplit partition. Peak RSS grew by about 393 MB, 
against 10 to 87 MB for Spark's reader. The optimized native library was much 
faster than Spark (217 ms against 673 ms warm), so the reader itself looks 
right and my concern is the memory and the up-front work.
   
   The comment in `CometScanRule` relies on the unsplit restriction to keep 
this to small lookup files, but that restriction only caps a file at 
`spark.sql.files.maxPartitionBytes`. The default is 128 MB and many jobs raise 
it. The docs describe the feature as targeting small lookup tables, and nothing 
enforces that. With several tasks per executor the unreserved memory multiplies.
   
   Would one of these work? Stream the object with `GetResult::into_stream()` 
and carry a partial line between chunks. Reserve the buffers against the task 
pool. Or enforce "small" in `CometScanRule` with an explicit per-file ceiling 
and fall back above it. The third is the least code. `LIMIT` also cannot stop 
early today because every batch is built up front.



##########
spark/src/test/scala/org/apache/comet/text/CometTextNativeReadSuite.scala:
##########
@@ -0,0 +1,279 @@
+/*
+ * 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.
+ */
+
+package org.apache.comet.text
+
+import java.io.File
+import java.nio.charset.StandardCharsets
+import java.nio.file.Files
+
+import org.apache.spark.sql.CometTestBase
+import org.apache.spark.sql.functions.input_file_name
+import org.apache.spark.sql.internal.SQLConf
+
+import org.apache.comet.CometConf
+
+class CometTextNativeReadSuite extends CometTestBase {

Review Comment:
   The PR is motivated by a text lookup on the build side of a broadcast join, 
but no test covers that case. I ran a Parquet fact table joined to a text 
lookup and the whole plan is native, so a `checkSparkAnswerAndOperator` test 
would lock in the point of the change. It should set 
`spark.sql.sources.useV1SourceList` to `avro,csv,json,kafka,orc,parquet`. The 
empty list that `withTextConf` uses also sends Parquet through the unsupported 
V2 scan.
   
   Two more tests would help. One reads several unsplit files that land in 
separate partitions, for example with a large 
`spark.sql.files.openCostInBytes`, since `file_partitions[self.partition]` is 
only exercised at index 0 today. The other covers the partition-column 
fallback, where selecting only `value` from a partitioned directory stays 
native and `select *` falls back.



##########
spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala:
##########
@@ -458,6 +459,147 @@ case class CometScanRule(session: SparkSession)
           withFallbackReasons(scanExec, fallbackReasons.toSet)
         }
 
+      case scan: TextScan if COMET_TEXT_V2_NATIVE_ENABLED.get() =>
+        if (scanExec.output.exists(_.isMetadataCol)) {
+          return withFallbackReason(
+            scanExec,
+            "Metadata columns are not supported for Text V2 scans")
+        }
+        // A native scan is useless without native execution: with exec 
disabled, CometExecRule
+        // leaves this CometBatchScanExec unconverted and it fails at runtime. 
Require native exec.
+        if (!COMET_EXEC_ENABLED.get()) {
+          return withFallbackReason(
+            scanExec,
+            s"Native Text scan requires ${COMET_EXEC_ENABLED.key} to be 
enabled")
+        }
+        // input_file_name/_block_start/_block_length read a thread-local that 
Spark's FileScanRDD
+        // sets; the native DataFusion scan does not, so they would return "" 
/ -1. Same guard as
+        // the native Parquet path.
+        if (plan.exists(node =>
+            node.expressions.exists(_.exists {
+              case _: InputFileName | _: InputFileBlockStart | _: 
InputFileBlockLength => true
+              case _ => false
+            }))) {
+          return withFallbackReason(
+            scanExec,
+            "Native Text scan is not compatible with input_file_name, " +
+              "input_file_block_start, or input_file_block_length")
+        }
+        val fallbackReasons = new ListBuffer[String]()
+        val schemaSupported =
+          CometBatchScanExec.isSchemaSupported(scan.readDataSchema, 
fallbackReasons)
+        if (!schemaSupported) {
+          fallbackReasons += s"Schema ${scan.readDataSchema} is not supported"
+        }
+        // Spark's text source allows only a single data column 
(TextScan.verifyReadSchema). A
+        // user-forced multi-column read schema is invalid; fall back so Spark 
raises its structured
+        // AnalysisException rather than letting the native reader fail with 
an opaque error.
+        if (scan.readDataSchema.size > 1) {
+          fallbackReasons += "Text data source supports only a single column"
+        }
+        // Partition columns are not implemented for native Text scans (unlike 
the native CSV path,
+        // which forwards them). Fall back to Spark for partitioned reads 
rather than drop them.
+        if (scan.readPartitionSchema.nonEmpty) {
+          fallbackReasons += "Comet does not support partition columns in 
native Text scans"
+        }
+        // The native reader loads each file whole. Spark splits a large 
uncompressed text file
+        // into byte-range partitions (TextScan.isSplitable), and each split 
would then re-read and
+        // buffer the whole file with no memory-pool reservation. Restrict 
native execution to
+        // unsplit files (the small build-side lookups this targets); Spark 
handles large files.
+        val hasSplitFile = scanExec.inputPartitions.exists {
+          case fp: FilePartition => fp.files.exists(_.start != 0)
+          case _ => false
+        }
+        if (hasSplitFile) {
+          fallbackReasons += "Comet native Text scan does not support files 
split into byte " +
+            "ranges (large files)"
+        }
+        // Comet's native plan does not carry a compression codec, so the 
native reader would feed
+        // raw compressed bytes to the line splitter. Spark auto-detects a 
codec by file extension
+        // (e.g. .gz, .bz2). Fall back when any input file is compressed. Use 
the options-aware
+        // Hadoop conf so a codec registered via a per-read option is honored 
(as Spark does).
+        val hadoopConf = 
session.sessionState.newHadoopConfWithOptions(scan.options.asScala.toMap)
+        val codecFactory = new 
org.apache.hadoop.io.compress.CompressionCodecFactory(hadoopConf)
+        val hasCompressedFile = scan.fileIndex.inputFiles.exists { path =>
+          codecFactory.getCodec(new org.apache.hadoop.fs.Path(path)) != null
+        }
+        if (hasCompressedFile) {
+          fallbackReasons += "Comet does not support compressed text files"
+        }
+        // Spark skips unreadable/missing files when these are set; the native 
reader's DataFusion
+        // FileStream fails the whole scan instead (OnError::Fail). Fall back 
to match Spark.
+        if (SQLConf.get.ignoreCorruptFiles || scan.options.getBoolean(
+            "ignoreCorruptFiles",
+            false)) {
+          fallbackReasons += "Comet native Text scan does not support 
ignoreCorruptFiles"
+        }
+        if (SQLConf.get.ignoreMissingFiles || scan.options.getBoolean(
+            "ignoreMissingFiles",
+            false)) {
+          fallbackReasons += "Comet native Text scan does not support 
ignoreMissingFiles"
+        }
+        // A wholetext file is never split (TextScan.isSplitable is false for 
wholeText), so the
+        // split-file guard above does not bound it. Reading a very large file 
whole is unbounded
+        // memory, and the whole file becomes one Arrow string value. Because 
a lossy UTF-8 decode
+        // can expand ill-formed bytes up to 3x, cap the raw file size at 
Int.MaxValue/3 so the
+        // decoded value stays under Arrow's 2GB 32-bit string-offset limit; 
fall back above that.
+        // (The native reader also fails cleanly rather than panicking if a 
value still overflows.)
+        val textOptions =
+          new org.apache.spark.sql.execution.datasources.text.TextOptions(
+            scan.options.asScala.toMap)
+        if (textOptions.wholeText) {
+          val maxWholeTextBytes =
+            math.min(session.sessionState.conf.filesMaxPartitionBytes, 
Int.MaxValue.toLong / 3)
+          val hasLargeWholeTextFile = scanExec.inputPartitions.exists {
+            case fp: FilePartition => fp.files.exists(_.length > 
maxWholeTextBytes)
+            case _ => false
+          }
+          if (hasLargeWholeTextFile) {
+            fallbackReasons += "Comet native Text scan does not support large 
wholetext files " +
+              s"(over ${SQLConf.FILES_MAX_PARTITION_BYTES.key}, capped below 
2GB)"
+          }
+        }
+        // Comet's native reader opens files through object_store, which only 
understands a fixed
+        // set of URL schemes. A text file on a custom Hadoop scheme 
(viewfs://, oss://, ...) would
+        // be claimed here then hard-fail at execution, whereas Spark reads it 
via the Hadoop FS
+        // API. Decline such schemes (same gate as the native Parquet path); 
hdfs:// routes through
+        // libhdfs and is left to native.
+        val libhdfsSchemes: Set[String] = COMET_LIBHDFS_SCHEMES.get() match {

Review Comment:
   The scheme and path checks from here down to `rejectedPath` are a copy of 
the Parquet V1 gate earlier in this file. The copy leaves out the 
`aliasScanBuckets` check. If I read `get_partitioned_files` in `planner.rs` 
right, native planning registers one object store per partition and strips the 
bucket from every key. A Text scan over `blob://a/x` and `blob://b/y` would 
then read both files from bucket `a`.
   
   The Parquet comment limits that guard to alias scans so it does not newly 
fall back plain multi-bucket scans that work by luck today. That reasoning does 
not apply to a scan that is new. Should the Text branch decline any 
multi-bucket scan? I would also pull the shared checks into one helper that 
both gates call.



##########
spark/src/main/scala/org/apache/spark/sql/comet/CometTextNativeScanExec.scala:
##########
@@ -0,0 +1,135 @@
+/*
+ * 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.
+ */
+
+package org.apache.spark.sql.comet
+
+import scala.jdk.CollectionConverters._
+
+import org.apache.spark.sql.catalyst.expressions.{Attribute, SortOrder}
+import org.apache.spark.sql.catalyst.plans.physical.{Partitioning, 
UnknownPartitioning}
+import org.apache.spark.sql.execution.SparkPlan
+import org.apache.spark.sql.execution.datasources.FilePartition
+import org.apache.spark.sql.execution.datasources.text.TextOptions
+import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
+import org.apache.spark.sql.execution.datasources.v2.text.TextScan
+
+import com.google.common.base.Objects
+
+import org.apache.comet.{CometConf, ConfigEntry}
+import org.apache.comet.objectstore.NativeConfig
+import org.apache.comet.serde.{CometOperatorSerde, OperatorOuterClass}
+import org.apache.comet.serde.OperatorOuterClass.Operator
+import org.apache.comet.serde.operator.{partition2Proto, schema2Proto}
+
+/*
+ * Native Text scan operator that delegates file reading to datafusion. 
Mirrors the native CSV
+ * scan path; produces a single `value: string` column.
+ */
+case class CometTextNativeScanExec(

Review Comment:
   Could this operator override `stringArgs` the way `CometNativeScanExec` and 
`CometIcebergNativeScanExec` do? The operator guide says to show the semantic 
fields and not the serialization state. Today the node prints the whole 
`nativeOp`, so `EXPLAIN`, the SQL tab and the event log carry the protobuf with 
every input file path. A lookup directory with thousands of files makes that 
large. `CometCsvNativeScanExec` has the same gap, so one change could fix both.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to