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]
