parthchandra commented on code in PR #6784:
URL: https://github.com/apache/datafusion-comet/pull/6784#discussion_r4235297151
##########
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:
Confirmed, fixed. Added a version-aware ShimFileFormat.isCompressedFile:
4.1/4.2 use HadoopCodecStreams (catches .gzip/.zstd), earlier versions use
CompressionCodecFactory. Added .gz and 4.1+ .gzip/.zstd fallback tests.
##########
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:
Done — now calls CometScanRule.readsInputFileBlock(plan).
##########
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:
Fixed. The scan now falls back when
mapreduce.input.linerecordreader.line.maxlength is set in the Hadoop conf (read
option, core-site.xml, or static spark.hadoop.*). Added a test.
##########
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:
Done — pulled the scheme/bucket/path checks into a shared
checkObjectStoreRootPaths helper that both the Parquet and Text gates call. The
Parquet gate keeps its existing messages and early-returns, so its behavior is
unchanged; the logic just lives in one place now. Also added the multi-bucket
alias decline to the Text gate.
##########
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:
Both fixed. withTextConf now sets
useV1SourceList=avro,csv,json,kafka,orc,parquet so Parquet stays native. Added
a broadcast-join test, a multi-file/separate-partitions test, and a
partition-column fallback test.
##########
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:
Fixed all three: the size-limit wording, the partition-column note
(selecting value stays native), and a link to the non-UTF-8 section.
##########
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:
Done. Added stringArgs = Iterator(output) to both CometTextNativeScanExec
and CometCsvNativeScanExec, so neither prints the full nativeOp (and its file
paths) in EXPLAIN.
##########
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:
Went with your option (c): a new spark.comet.scan.text.maxFileSize (default
64MB), and the scan falls back above it. That caps the memory regardless of
maxPartitionBytes. Streaming the read is a good follow-up if we want to remove
the limit later.
--
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]