diff --git a/file/src/main/java/org/apache/pulsar/io/file/utils/ZipFiles.java b/file/src/main/java/org/apache/pulsar/io/file/utils/ZipFiles.java index bc4cfc3ffb..13e729f1df 100644 --- a/file/src/main/java/org/apache/pulsar/io/file/utils/ZipFiles.java +++ b/file/src/main/java/org/apache/pulsar/io/file/utils/ZipFiles.java @@ -29,7 +29,13 @@ import java.io.UncheckedIOException; import java.nio.file.Files; import java.nio.file.Path; +import java.util.Iterator; +import java.util.NoSuchElementException; +import java.util.Spliterator; +import java.util.Spliterators; import java.util.stream.Stream; +import java.util.stream.StreamSupport; +import java.util.zip.ZipEntry; import java.util.zip.ZipInputStream; /** @@ -51,34 +57,102 @@ public static boolean isZip(File f) { } /** - * Get a lazily loaded stream of lines from a gzipped file, similar to + * Get a lazily loaded stream of lines from every file entry of a zip file, similar to * {@link Files#lines(java.nio.file.Path)}. * + *

A {@link ZipInputStream} yields no data until it is positioned onto an entry via + * {@link ZipInputStream#getNextEntry()}; without that call the reader sees an empty stream + * and no lines are produced. The returned stream lazily walks every file entry and emits + * their lines in entry order, so multi-entry archives contribute all of their lines without + * buffering the whole archive in memory. The caller must close the returned stream to + * release the underlying file. + * * @param path * The path to the zipped file. - * @return stream with lines. + * @return stream with the lines of all file entries, in entry order. */ public static Stream lines(Path path) { - ZipInputStream zipStream = null; - + ZipInputStream zipStream; try { - zipStream = new ZipInputStream(Files.newInputStream(path)); + zipStream = new ZipInputStream(Files.newInputStream(path)); } catch (IOException e) { - closeSafely(zipStream); - throw new UncheckedIOException(e); + throw new UncheckedIOException(e); + } + Spliterator spliterator = Spliterators.spliteratorUnknownSize( + new ZipLineIterator(zipStream), Spliterator.ORDERED | Spliterator.NONNULL); + return StreamSupport.stream(spliterator, false).onClose(() -> closeSafely(zipStream)); + } + + /** + * Iterates the lines of every non-directory entry of a zip stream in entry order, advancing + * to the next entry when the current one is exhausted. Entry boundaries are line boundaries: + * a fresh {@link BufferedReader} is used for each entry, so an entry whose content does not + * end with a newline does not merge its trailing text into the next entry's first line. + */ + private static final class ZipLineIterator implements Iterator { + private final ZipInputStream zipStream; + private BufferedReader reader; + private String nextLine; + + ZipLineIterator(ZipInputStream zipStream) { + this.zipStream = zipStream; + } + + @Override + public boolean hasNext() { + if (nextLine == null) { + nextLine = readNextLine(); + } + return nextLine != null; + } + + @Override + public String next() { + if (!hasNext()) { + throw new NoSuchElementException(); + } + String line = nextLine; + nextLine = null; + return line; + } + + private String readNextLine() { + try { + while (true) { + if (reader != null) { + String line = reader.readLine(); + if (line != null) { + return line; + } + // Current entry is exhausted; drop its reader but keep the shared + // zipStream open so we can position onto the next entry. + reader = null; + } + ZipEntry entry = zipStream.getNextEntry(); + if (entry == null) { + return null; + } + if (!entry.isDirectory()) { + // Wrap the shared zipStream in a fresh reader for this entry. read() + // returns -1 at the entry boundary, so the reader stops before the next + // entry; the reader is intentionally not closed, as that would close the + // shared zipStream mid-iteration. + reader = new BufferedReader(new InputStreamReader(zipStream)); + } + } + } catch (IOException e) { + throw new UncheckedIOException(e); + } } - // Reader decoder = new InputStreamReader(gzipStream, Charset.defaultCharset()); - BufferedReader reader = new BufferedReader(new InputStreamReader(zipStream)); - return reader.lines().onClose(() -> closeSafely(reader)); } private static void closeSafely(Closeable closeable) { if (closeable != null) { - try { - closeable.close(); - } catch (IOException e) { - // Ignore - } + try { + closeable.close(); + } catch (IOException e) { + // Ignore + } } } } diff --git a/file/src/test/java/org/apache/pulsar/io/file/utils/ZipFilesTest.java b/file/src/test/java/org/apache/pulsar/io/file/utils/ZipFilesTest.java index 64f3ff42c7..99f85f1f7e 100644 --- a/file/src/test/java/org/apache/pulsar/io/file/utils/ZipFilesTest.java +++ b/file/src/test/java/org/apache/pulsar/io/file/utils/ZipFilesTest.java @@ -18,12 +18,19 @@ */ package org.apache.pulsar.io.file.utils; +import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; import static org.testng.Assert.assertTrue; import java.io.File; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; +import java.util.List; +import java.util.stream.Collectors; import java.util.stream.Stream; +import java.util.zip.ZipEntry; +import java.util.zip.ZipOutputStream; import org.testng.annotations.Test; public class ZipFilesTest { @@ -52,10 +59,38 @@ public final void nonExistantGzipFileTest() { public final void streamZipFileTest() { Path path = Paths.get(getFile("org/apache/pulsar/io/file/validZip.zip").getAbsolutePath(), ""); + // validZip.zip contains a single entry with the nine lines "Line 1".."Line 9". try (Stream lines = ZipFiles.lines(path)) { - lines.forEachOrdered(line -> assertTrue(line.startsWith("Line "))); - } catch (Exception e) { - e.printStackTrace(); + List collected = lines.collect(Collectors.toList()); + assertEquals(collected.size(), 9, "expected nine lines from the zip entry"); + for (int i = 0; i < collected.size(); i++) { + assertEquals(collected.get(i), "Line " + (i + 1)); + } + } + } + + @Test + public final void streamMultiEntryZipFileTest() throws Exception { + // Build a two-entry archive in a temp file so the test is self-contained and proves + // that lines from every entry are returned, in entry order. + Path zip = Files.createTempFile("pulsar-io-file-ziptest", ".zip"); + try { + try (ZipOutputStream out = new ZipOutputStream(Files.newOutputStream(zip))) { + out.putNextEntry(new ZipEntry("first.txt")); + out.write("a1\na2".getBytes(StandardCharsets.UTF_8)); + out.closeEntry(); + out.putNextEntry(new ZipEntry("second.txt")); + out.write("b1\nb2\nb3".getBytes(StandardCharsets.UTF_8)); + out.closeEntry(); + } + + assertTrue(ZipFiles.isZip(zip.toFile())); + try (Stream lines = ZipFiles.lines(zip)) { + assertEquals(lines.collect(Collectors.toList()), + List.of("a1", "a2", "b1", "b2", "b3")); + } + } finally { + Files.deleteIfExists(zip); } }