This is an automated email from the ASF dual-hosted git repository.

cloud-fan pushed a commit to branch branch-4.x
in repository https://gitbox.apache.org/repos/asf/spark.git


The following commit(s) were added to refs/heads/branch-4.x by this push:
     new 6ce431bbfe24 [SPARK-57478][SQL] Read text files from tar archives
6ce431bbfe24 is described below

commit 6ce431bbfe242019835676f0f8d6a703faad2078
Author: akshatshenoi-db <[email protected]>
AuthorDate: Wed Jun 17 10:31:54 2026 -0700

    [SPARK-57478][SQL] Read text files from tar archives
    
    ### What changes were proposed in this pull request?
    
    SPARK-57135 added reading CSV files packed in tar archives 
(`.tar`/`.tar.gz`/`.tgz`) and SPARK-57321 added schema inference for them; 
SPARK-57419 extended both to JSON. All are gated by 
`spark.sql.files.archive.reader.enabled`. This extends archive reading to the 
text data source.
    
    When the flag is enabled, the V1 text data source reads a tar archive as if 
it were a directory of its entries: each entry is streamed through 
`ArchiveReader` (never unpacked to disk) and read exactly like a standalone 
text file -- one row per line, or a single row holding the whole entry when 
`wholeText` is set (`TextFileFormat.readArchive`). The whole archive is one 
non-splittable unit (`isSplitable` returns false for an archive path), and a 
corrupt/missing archive is skipped as a u [...]
    
    Text has a fixed `value STRING` schema, so there is no schema inference. 
Archive scanning is wired into the V1 file source only; the DSv2 reader is left 
untouched.
    
    ### Why are the changes needed?
    
    To let text ingestion read tar archives without unpacking them to disk, 
matching the CSV and JSON behavior already in Spark.
    
    ### Does this PR introduce _any_ user-facing change?
    
    Yes. With `spark.sql.files.archive.reader.enabled=true` (default false), 
the text data source can read `.tar`/`.tar.gz`/`.tgz` files.
    
    ### How was this patch tested?
    
    New `TextTarArchiveReadSuite`: reads of multi-entry archives across all 
three extensions, parity with a directory read of the same files, `wholeText` 
and a custom line separator, empty and corrupt archives, the single-partition 
guarantee, and an archive mixed with loose files in the same directory.
    
    ### Was this patch authored or co-authored using generative AI tooling?
    
    Generated-by: Claude Code
    
    Closes #56527 from akshatshenoi-db/archive-text.
    
    Authored-by: akshatshenoi-db <[email protected]>
    Signed-off-by: Wenchen Fan <[email protected]>
    (cherry picked from commit 7eec4f5b657fa2b81cbc7391b49e31fef3289944)
    Signed-off-by: Wenchen Fan <[email protected]>
---
 .../org/apache/spark/sql/internal/SQLConf.scala    |   9 +-
 .../datasources/text/TextFileFormat.scala          |  58 ++++++-
 .../execution/datasources/TarArchiveReadBase.scala |  41 +----
 ...iveReadBase.scala => TarArchiveTestUtils.scala} |  21 +--
 .../datasources/TextTarArchiveReadSuite.scala      | 178 +++++++++++++++++++++
 5 files changed, 256 insertions(+), 51 deletions(-)

diff --git 
a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala 
b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala
index 0a94c5a5980a..89f606b5ef6a 100644
--- a/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala
+++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala
@@ -2726,10 +2726,11 @@ object SQLConf {
     .createWithDefaultString("128MB") // parquet.block.size
 
   val ARCHIVE_FORMAT_READER_ENABLED = 
buildConf("spark.sql.files.archive.reader.enabled")
-    .doc("When true, the CSV data source can read tar archives (.tar, .tar.gz, 
.tgz): each " +
-      "archive is read as a single split and its entries are streamed through 
the CSV parser " +
-      "(never unpacked to disk), as if the entries were separate CSV files, 
both during scan " +
-      "and schema inference. Only the CSV data source supports reading 
archives.")
+    .doc("When true, a supported data source can read tar archives (.tar, 
.tar.gz, .tgz): " +
+      "each archive is read as a single split and its entries are streamed 
through that data " +
+      "source's parser (never unpacked to disk), as if the entries were 
separate files, both " +
+      "during scan and schema inference. The CSV, JSON, and text data sources 
support " +
+      "reading archives.")
     .version("5.0.0")
     .withBindingPolicy(ConfigBindingPolicy.SESSION)
     .booleanConf
diff --git 
a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/text/TextFileFormat.scala
 
b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/text/TextFileFormat.scala
index 412912fbe6c5..863f4e1c162e 100644
--- 
a/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/text/TextFileFormat.scala
+++ 
b/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/text/TextFileFormat.scala
@@ -60,6 +60,10 @@ case class TextFileFormat() extends TextBasedFileFormat with 
DataSourceRegister
       options: Map[String, String],
       path: Path): Boolean = {
     val textOptions = new TextOptions(options)
+    if (textOptions.archiveFormatEnabled && ArchiveReader.isArchivePath(path)) 
{
+      // A tar archive is read as one sequential stream (entry by entry), so 
it is never split.
+      return false
+    }
     super.isSplitable(sparkSession, options, path) && !textOptions.wholeText
   }
 
@@ -108,7 +112,59 @@ case class TextFileFormat() extends TextBasedFileFormat 
with DataSourceRegister
     val textOptions = new TextOptions(options)
     val broadcastedHadoopConf =
       SerializableConfiguration.broadcast(sparkSession.sparkContext, 
hadoopConf)
-    readToUnsafeMem(broadcastedHadoopConf, requiredSchema, textOptions)
+    val perFileReader = readToUnsafeMem(broadcastedHadoopConf, requiredSchema, 
textOptions)
+    val archiveReader = readArchive(broadcastedHadoopConf, requiredSchema, 
textOptions)
+    // A tar archive (always a single split, see `isSplitable`) is streamed 
entry by entry when
+    // archive reads are enabled; otherwise the file is read directly. Archive 
scanning is wired
+    // into the V1 file source only, so this dispatch lives here rather than 
in a shared reader.
+    (file: PartitionedFile) => {
+      if (textOptions.archiveFormatEnabled && 
ArchiveReader.isArchivePath(file.toPath)) {
+        archiveReader(file)
+      } else {
+        perFileReader(file)
+      }
+    }
+  }
+
+  /**
+   * Streams a tar archive (`.tar`/`.tar.gz`/`.tgz`) entry by entry, emitting 
the same
+   * `value`-column rows the per-file reader produces -- each entry is read as 
if it were a
+   * standalone text file (one row per line, or a single row holding the whole 
entry when
+   * `wholeText` is set), without unpacking the archive to disk. The whole 
archive is a single
+   * split (see `isSplitable`).
+   */
+  private def readArchive(
+      conf: Broadcast[SerializableConfiguration],
+      requiredSchema: StructType,
+      textOptions: TextOptions): PartitionedFile => Iterator[UnsafeRow] = {
+    (file: PartitionedFile) => {
+      val confValue = conf.value.value
+      ArchiveReader(file.toPath).readEntries(confValue) { (_, in) =>
+        // Each entry is read as a standalone text file, so it gets its own 
row writer, exactly as
+        // `readToUnsafeMem` builds one per file.
+        val emptyUnsafeRow = new UnsafeRow(0)
+        val unsafeRowWriter = new UnsafeRowWriter(1)
+        // Mirrors `readToUnsafeMem`: an empty required schema (e.g. `count`) 
yields one empty row
+        // per record; otherwise each record is written into the single 
`value` column.
+        def toRow(bytes: Array[Byte], length: Int): UnsafeRow = {
+          if (requiredSchema.isEmpty) {
+            emptyUnsafeRow
+          } else {
+            unsafeRowWriter.reset()
+            unsafeRowWriter.write(0, bytes, 0, length)
+            unsafeRowWriter.getRow()
+          }
+        }
+        if (textOptions.wholeText) {
+          val content = in.readAllBytes()
+          Iterator.single(toRow(content, content.length))
+        } else {
+          ArchiveReader.lineIterator(in, textOptions.lineSeparatorInRead).map 
{ line =>
+            toRow(line.getBytes, line.getLength)
+          }
+        }
+      }
+    }
   }
 
   private def readToUnsafeMem(
diff --git 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
index 60a73b0891ba..234bd7a2f1f2 100644
--- 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
+++ 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
@@ -17,46 +17,13 @@
 
 package org.apache.spark.sql.execution.datasources
 
-import java.io.{File, FileOutputStream, OutputStream}
-import java.nio.charset.StandardCharsets
-import java.nio.file.Files
-import java.util.Locale
-import java.util.zip.GZIPOutputStream
-
-import org.apache.commons.compress.archivers.tar.{TarArchiveEntry, 
TarArchiveOutputStream}
-
 /**
  * Binds [[ArchiveReadSuiteBase]]'s archive-format hooks to tar containers: 
plain `.tar`, gzipped
- * `.tar.gz`, and `.tgz`. Reusable across file formats -- a 
`<format>TarArchiveReadSuite` mixes this
- * in alongside the file-format trait.
+ * `.tar.gz`, and `.tgz`. The container-writing helpers live in 
[[TarArchiveTestUtils]] (shared with
+ * standalone suites that cannot extend `ArchiveReadSuiteBase`). Reusable 
across file formats -- a
+ * `<format>TarArchiveReadSuite` mixes this in alongside the file-format trait.
  */
-trait TarArchiveReadBase extends ArchiveReadSuiteBase {
-
-  override protected def archiveExtensions: Seq[String] = Seq("tar", "tar.gz", 
"tgz")
+trait TarArchiveReadBase extends ArchiveReadSuiteBase with TarArchiveTestUtils 
{
 
   override protected def corruptArchiveExtension: String = "tar.gz"
-
-  override protected def writeArchive(dest: File, entries: Seq[(String, 
Array[Byte])]): Unit = {
-    val name = dest.getName.toLowerCase(Locale.ROOT)
-    val rawOut: OutputStream = if (name.endsWith(".gz") || 
name.endsWith(".tgz")) {
-      new GZIPOutputStream(new FileOutputStream(dest))
-    } else {
-      new FileOutputStream(dest)
-    }
-    val out = new TarArchiveOutputStream(rawOut)
-    try {
-      entries.foreach { case (entryName, bytes) =>
-        val entry = new TarArchiveEntry(entryName)
-        entry.setSize(bytes.length.toLong)
-        out.putArchiveEntry(entry)
-        out.write(bytes)
-        out.closeArchiveEntry()
-      }
-      out.finish()
-    } finally out.close()
-  }
-
-  override protected def writeCorruptArchive(dest: File): Unit =
-    Files.write(dest.toPath, "this is not a valid gzip-compressed tar archive"
-      .getBytes(StandardCharsets.UTF_8))
 }
diff --git 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveTestUtils.scala
similarity index 66%
copy from 
sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
copy to 
sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveTestUtils.scala
index 60a73b0891ba..7f1d48e8143b 100644
--- 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveReadBase.scala
+++ 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TarArchiveTestUtils.scala
@@ -26,17 +26,19 @@ import java.util.zip.GZIPOutputStream
 import org.apache.commons.compress.archivers.tar.{TarArchiveEntry, 
TarArchiveOutputStream}
 
 /**
- * Binds [[ArchiveReadSuiteBase]]'s archive-format hooks to tar containers: 
plain `.tar`, gzipped
- * `.tar.gz`, and `.tgz`. Reusable across file formats -- a 
`<format>TarArchiveReadSuite` mixes this
- * in alongside the file-format trait.
+ * Tar container-writing helpers (plain `.tar`, gzipped `.tar.gz`, and 
`.tgz`), independent of any
+ * read-test harness. [[TarArchiveReadBase]] mixes this into 
[[ArchiveReadSuiteBase]] for the
+ * format-agnostic suites, and standalone suites that cannot extend 
`ArchiveReadSuiteBase` (e.g.
+ * `TextTarArchiveReadSuite`, whose single `value` column doesn't fit the 
two-column shared tests)
+ * mix it in directly, so the container logic lives in one place.
  */
-trait TarArchiveReadBase extends ArchiveReadSuiteBase {
+trait TarArchiveTestUtils {
 
-  override protected def archiveExtensions: Seq[String] = Seq("tar", "tar.gz", 
"tgz")
+  /** Tar extensions to exercise; the head is the default. */
+  protected def archiveExtensions: Seq[String] = Seq("tar", "tar.gz", "tgz")
 
-  override protected def corruptArchiveExtension: String = "tar.gz"
-
-  override protected def writeArchive(dest: File, entries: Seq[(String, 
Array[Byte])]): Unit = {
+  /** Writes `entries` (name -> bytes) into the archive at `dest`; compression 
follows the ext. */
+  protected def writeArchive(dest: File, entries: Seq[(String, Array[Byte])]): 
Unit = {
     val name = dest.getName.toLowerCase(Locale.ROOT)
     val rawOut: OutputStream = if (name.endsWith(".gz") || 
name.endsWith(".tgz")) {
       new GZIPOutputStream(new FileOutputStream(dest))
@@ -56,7 +58,8 @@ trait TarArchiveReadBase extends ArchiveReadSuiteBase {
     } finally out.close()
   }
 
-  override protected def writeCorruptArchive(dest: File): Unit =
+  /** Writes bytes that are not a valid gzip-compressed tar archive to `dest`. 
*/
+  protected def writeCorruptArchive(dest: File): Unit =
     Files.write(dest.toPath, "this is not a valid gzip-compressed tar archive"
       .getBytes(StandardCharsets.UTF_8))
 }
diff --git 
a/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TextTarArchiveReadSuite.scala
 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TextTarArchiveReadSuite.scala
new file mode 100644
index 000000000000..0163d03a168a
--- /dev/null
+++ 
b/sql/core/src/test/scala/org/apache/spark/sql/execution/datasources/TextTarArchiveReadSuite.scala
@@ -0,0 +1,178 @@
+/*
+ * 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.execution.datasources
+
+import java.io.File
+import java.nio.charset.StandardCharsets
+import java.nio.file.Files
+
+import org.apache.spark.{SparkConf, SparkException}
+import org.apache.spark.sql.{DataFrame, QueryTest, Row}
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.test.SharedSparkSession
+import org.apache.spark.util.Utils
+
+/**
+ * Reads of text files packed in tar archives (`.tar`/`.tar.gz`/`.tgz`), 
streamed through the
+ * [[ArchiveReader]] path. Entries are streamed (never unpacked to disk), and 
the central contract
+ * verified throughout is parity with reading the same files from a directory.
+ *
+ * Unlike CSV/JSON this does not reuse [[ArchiveReadSuiteBase]]: the text data 
source has a single
+ * fixed `value` column (one row per line, or per entry with `wholetext`) and 
no schema inference,
+ * so the structured, two-column tests there do not apply. The tar 
container-writing helpers come
+ * from [[TarArchiveTestUtils]] (shared with [[TarArchiveReadBase]]).
+ */
+class TextTarArchiveReadSuite extends QueryTest with SharedSparkSession with 
TarArchiveTestUtils {
+
+  override def sparkConf: SparkConf =
+    super.sparkConf.set(SQLConf.ARCHIVE_FORMAT_READER_ENABLED.key, "true")
+
+  private def textBytes(s: String): Array[Byte] = 
s.getBytes(StandardCharsets.UTF_8)
+
+  /** Provides an archive-extensioned path inside a fresh temp dir to `f`. */
+  private def withArchiveFile(
+      extension: String = archiveExtensions.head)(f: File => Unit): Unit = {
+    val dir = Utils.createTempDir(namePrefix = "archive-test")
+    try f(new File(dir, s"archive.$extension")) finally 
Utils.deleteRecursively(dir)
+  }
+
+  private def read(path: String, options: Map[String, String] = Map.empty): 
DataFrame =
+    spark.read.options(options).text(path)
+
+  test("read a tar archive of multiple text entries matches the union of the 
lines") {
+    archiveExtensions.foreach { ext =>
+      withArchiveFile(ext) { archive =>
+        writeArchive(archive, Seq(
+          "a.txt" -> textBytes("line1\nline2\n"),
+          "b.txt" -> textBytes("line3\n"),
+          "c.txt" -> textBytes("line4\nline5\n")))
+        checkAnswer(
+          read(archive.getCanonicalPath),
+          Seq("line1", "line2", "line3", "line4", "line5").map(Row(_)))
+      }
+    }
+  }
+
+  test("archive entries read like a directory of the same files") {
+    val entries = Seq("a.txt" -> textBytes("a1\na2\n"), "b.txt" -> 
textBytes("b1\n"))
+    withArchiveFile() { archive =>
+      writeArchive(archive, entries)
+      val fromArchive = read(archive.getCanonicalPath)
+      withTempDir { dir =>
+        entries.foreach { case (n, b) => Files.write(new File(dir, n).toPath, 
b) }
+        checkAnswer(fromArchive, read(dir.getCanonicalPath).collect().toSeq)
+      }
+    }
+  }
+
+  test("an empty archive yields no rows") {
+    withArchiveFile() { archive =>
+      writeArchive(archive, Seq.empty)
+      checkAnswer(read(archive.getCanonicalPath), Seq.empty[Row])
+    }
+  }
+
+  test("an archive and loose text files in the same directory are all read") {
+    withTempDir { dir =>
+      val ext = archiveExtensions.head
+      writeArchive(
+        new File(dir, s"data.$ext"),
+        Seq("a.txt" -> textBytes("in-archive-1\nin-archive-2\n")))
+      Files.write(new File(dir, "loose.txt").toPath, textBytes("loose-1\n"))
+      checkAnswer(
+        read(dir.getCanonicalPath),
+        Seq("in-archive-1", "in-archive-2", "loose-1").map(Row(_)))
+    }
+  }
+
+  test("wholetext reads each entry as a single row") {
+    withArchiveFile() { archive =>
+      writeArchive(archive, Seq(
+        "a.txt" -> textBytes("l1\nl2"),
+        "b.txt" -> textBytes("only")))
+      checkAnswer(
+        read(archive.getCanonicalPath, Map("wholetext" -> "true")),
+        Seq(Row("l1\nl2"), Row("only")))
+    }
+  }
+
+  test("a custom line separator splits entries into rows") {
+    withArchiveFile() { archive =>
+      writeArchive(archive, Seq("a.txt" -> textBytes("x;y;z")))
+      checkAnswer(
+        read(archive.getCanonicalPath, Map("lineSep" -> ";")),
+        Seq(Row("x"), Row("y"), Row("z")))
+    }
+  }
+
+  test("count over an archive reads the right number of rows with an empty 
required schema") {
+    withArchiveFile() { archive =>
+      writeArchive(archive, Seq(
+        "a.txt" -> textBytes("1\n2\n3\n"),
+        "b.txt" -> textBytes("4\n")))
+      assert(read(archive.getCanonicalPath).count() == 4L)
+    }
+  }
+
+  test("an archive always yields a single partition regardless of size") {
+    withArchiveFile() { archive =>
+      val big = (1 to 1000).map(i => s"value-$i").mkString("\n")
+      writeArchive(archive, (0 until 4).map(i => s"part-$i.txt" -> 
textBytes(big + "\n")))
+      withSQLConf(SQLConf.FILES_MAX_PARTITION_BYTES.key -> "1024") {
+        val readDf = read(archive.getCanonicalPath)
+        assert(readDf.rdd.getNumPartitions == 1,
+          s"archive should be a single partition; got 
${readDf.rdd.getNumPartitions}")
+        assert(readDf.count() == 4000L)
+      }
+    }
+  }
+
+  Seq(true, false).foreach { ignoreCorrupt =>
+    test(s"ignoreCorruptFiles=$ignoreCorrupt controls whether a corrupt 
archive is skipped") {
+      withArchiveFile("tar.gz") { archive =>
+        writeCorruptArchive(archive)
+        withSQLConf(SQLConf.IGNORE_CORRUPT_FILES.key -> 
ignoreCorrupt.toString) {
+          if (ignoreCorrupt) {
+            checkAnswer(read(archive.getCanonicalPath), Seq.empty[Row])
+          } else {
+            intercept[SparkException](read(archive.getCanonicalPath).collect())
+          }
+        }
+      }
+    }
+  }
+
+  Seq(true, false).foreach { ignoreMissing =>
+    test(s"ignoreMissingFiles=$ignoreMissing controls whether a missing 
archive is skipped") {
+      withArchiveFile() { archive =>
+        writeArchive(archive, Seq("a.txt" -> textBytes("line1\nline2\n")))
+        withSQLConf(SQLConf.IGNORE_MISSING_FILES.key -> 
ignoreMissing.toString) {
+          // The archive is listed when the DataFrame is built, then deleted 
before the scan opens
+          // it, so the reader hits a missing file -- handled by 
`FileScanRDD`, like any file.
+          val df = read(archive.getCanonicalPath)
+          assert(archive.delete(), s"failed to delete $archive")
+          if (ignoreMissing) {
+            checkAnswer(df, Seq.empty[Row])
+          } else {
+            intercept[SparkException](df.collect())
+          }
+        }
+      }
+    }
+  }
+}


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

Reply via email to