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


##########
spark/src/main/scala/org/apache/spark/sql/comet/CometTextNativeScanExec.scala:
##########
@@ -0,0 +1,140 @@
+/*
+ * 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(
+    override val nativeOp: Operator,
+    override val output: Seq[Attribute],
+    @transient override val originalPlan: BatchScanExec,
+    override val serializedPlanOpt: SerializedPlan)
+    extends CometLeafExec {
+  override val supportsColumnar: Boolean = true
+
+  override val nodeName: String = "CometTextNativeScan"
+
+  // Show only the semantic output in EXPLAIN / the SQL tab / the event log. 
The default would print
+  // the whole `nativeOp` protobuf, which carries every input file path (large 
for a lookup
+  // directory of many files). Mirrors CometNativeScanExec.
+  override def stringArgs: Iterator[Any] = Iterator(output)
+
+  override def outputPartitioning: Partitioning = UnknownPartitioning(
+    originalPlan.inputPartitions.length)
+
+  override def outputOrdering: Seq[SortOrder] = Nil
+
+  override protected def doCanonicalize(): SparkPlan = {
+    CometTextNativeScanExec(nativeOp, output, originalPlan, serializedPlanOpt)
+  }
+
+  override def equals(obj: Any): Boolean = {
+    obj match {
+      case other: CometTextNativeScanExec =>
+        this.output == other.output &&
+        this.serializedPlanOpt == other.serializedPlanOpt &&
+        this.originalPlan == other.originalPlan
+      case _ =>
+        false
+    }
+  }
+
+  override def hashCode(): Int = {
+    Objects.hashCode(output, serializedPlanOpt, originalPlan)
+  }
+}
+
+object CometTextNativeScanExec extends CometOperatorSerde[CometBatchScanExec] {
+
+  override def enabledConfig: Option[ConfigEntry[Boolean]] = Some(
+    CometConf.COMET_TEXT_V2_NATIVE_ENABLED)
+
+  override def convert(
+      op: CometBatchScanExec,
+      builder: Operator.Builder,
+      childOp: Operator*): Option[Operator] = {
+    val textScanBuilder = OperatorOuterClass.TextScan.newBuilder()
+    val textScan = op.wrapped.scan.asInstanceOf[TextScan]
+    val sessionState = op.session.sessionState
+    val options = new TextOptions(textScan.options.asScala.toMap)
+    val filePartitions = op.inputPartitions.map(_.asInstanceOf[FilePartition])
+    val textOptionsProto = textOptions2Proto(options)
+    val dataSchemaProto = schema2Proto(textScan.dataSchema)
+    val readSchemaFieldNames = textScan.readDataSchema.fieldNames
+    val projectionVector = textScan.dataSchema.zipWithIndex
+      .filter { case (field, _) =>
+        readSchemaFieldNames.contains(field.name)
+      }
+      .map(_._2.asInstanceOf[Integer])
+    val partitionSchemaProto = schema2Proto(textScan.readPartitionSchema)
+    val partitionsProto = filePartitions.map(partition2Proto(_, 
textScan.readPartitionSchema))
+
+    val objectStoreOptions = filePartitions.headOption
+      .flatMap { partitionFile =>
+        val hadoopConf = sessionState
+          .newHadoopConfWithOptions(op.session.sparkContext.conf.getAll.toMap)

Review Comment:
   [P2] Preserve per-read Hadoop settings when serializing the text scan. For 
example, a text read using `.option("fs.s3a.endpoint", 
"http://127.0.0.1:19000";)` expects that endpoint to be used. Spark's reader and 
the new eligibility gate incorporate read options, but this call rebuilds the 
configuration from `SparkConf` and discards them. `extractObjectStoreOptions` 
therefore sends the global/default endpoint to native execution, so enabling 
native text can break an otherwise working read against an S3-compatible 
service. Build this configuration from 
`textScan.options.asCaseSensitiveMap.asScala.toMap`, matching Spark's 
`TextScan.createReaderFactory`.
   
   Evidence: A bounded Spark 4.1.3 JVM probe executed both configuration 
constructions with the endpoint supplied only in read options. Spark's 
options-aware configuration returned `http://127.0.0.1:19000`; the exact 
construction used at line 112 returned null. 
`NativeConfig.extractObjectStoreOptions` copies `fs.s3a.*` exclusively from 
that supplied configuration, and the native text planner passes the resulting 
map to object-store creation. No remote service or credentials were used in 
this reproduction.



##########
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?;
+
+            // Decompress only when the file is actually compressed. On 
Comet's plan config the file
+            // is always UNCOMPRESSED (CometScanRule falls back for compressed 
files), so we borrow
+            // the fetched bytes directly and avoid a full-file copy; the 
decompress arm is a
+            // defensive handler should compression ever be plumbed through. A 
byte range never
+            // coexists with decompression (compressed files are 
non-splittable).
+            //
+            // NOTE: the whole object is read into memory. CometScanRule 
restricts native text to
+            // unsplit files (it falls back when Spark splits a file into byte 
ranges), so in
+            // practice this reads each small file once; the range filter 
below is a correctness
+            // safety net for any split that still reaches here.
+            let buf: Cow<[u8]> = if file_compression_type.is_compressed() {
+                let mut decoder = 
file_compression_type.convert_read(std::io::Cursor::new(raw))?;
+                let mut decoded = Vec::new();
+                decoder.read_to_end(&mut decoded)?;
+                Cow::Owned(decoded)
+            } else {
+                Cow::Borrowed(raw.as_ref())
+            };
+
+            let mut timer = baseline_metrics.elapsed_compute().timer();
+            // Spark's line reader (Hadoop 
LineRecordReader.skipUtfByteOrderMark) strips a leading
+            // UTF-8 BOM at the start of the file. Mirror that in line mode 
for the split that starts
+            // at file offset 0 (native only reads unsplit files, so start is 
always 0; a custom
+            // lineSep still strips, matching Hadoop). wholetext keeps the 
BOM, as Spark's
+            // WholeTextFileRecordReader does.
+            let at_file_start = range.as_ref().map(|r| r.start == 
0).unwrap_or(true);
+            let content = strip_leading_bom(&buf, whole_text, at_file_start);

Review Comment:
   [P2] Apply BOM handling after splitting the first record. With native V2 
text enabled, a UTF-8 file containing `\uFEFFa\uFEFFb` and `.option("lineSep", 
"\uFEFF")` should return `["", "a", "b"]`, as Spark does. This code removes the 
initial separator before splitting and returns `["a", "b"]`. With separator 
`\uFEFFx`, it also changes the first value from `a` to `xa`. These inputs pass 
the eligibility checks and silently change counts, filters, and join results. 
Hadoop reads the first record before checking its value for a BOM. Follow its 
first-record BOM/consumed-byte rules, or fall back for conflicting separators.
   
   Evidence: Compiled `strip_leading_bom`, `split_lines`, and their split 
helpers directly from the reviewed head. The Rust probe returned two rows for 
both examples. Hadoop `LineRecordReader`, Spark `HadoopLineRecordReader`, and 
Spark 4.1.3 V2 SQL returned three rows: ["", "a", "b"]. A file containing only 
U+FEFF with that separator returned zero native-helper rows versus one empty 
Spark row. Ordinary BOM and BOM-only default-separator controls agreed.



-- 
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