From 34c18ee517782d01c65829db9f4903adeca5d045 Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Sun, 9 Aug 2026 18:53:36 +0100 Subject: [PATCH 1/2] file: add compression level --- .../archive/FileByteStringSeparators.scala | 4 ++-- .../file/impl/archive/ZipArchiveFlow.scala | 12 +++++++--- .../file/impl/archive/ZipArchiveManager.scala | 5 +++-- .../connectors/file/javadsl/Archive.scala | 11 ++++++++-- .../connectors/file/scaladsl/Archive.scala | 9 +++++++- .../scala/docs/scaladsl/ArchiveSpec.scala | 14 ++++++++++++ .../impl/archive/ZipArchiveFlowTest.scala | 22 +++++++++++++++++++ 7 files changed, 67 insertions(+), 10 deletions(-) diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/FileByteStringSeparators.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/FileByteStringSeparators.scala index b4cd02be0..bc6c63082 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/FileByteStringSeparators.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/FileByteStringSeparators.scala @@ -35,10 +35,10 @@ import pekko.util.ByteString ByteString(endFileWord) def isStartingByteString(b: ByteString): Boolean = - b.utf8String.startsWith(startFileWord) + b.size >= 7 && b.slice(0, 7).utf8String == startFileWord def isEndingByteString(b: ByteString): Boolean = - b.utf8String == endFileWord + b.size == 5 && b.utf8String == endFileWord def getPathFromStartingByteString(b: ByteString): String = { val splitted = b.utf8String.split(separator) diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlow.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlow.scala index 41abf75ad..e5e99930c 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlow.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlow.scala @@ -26,13 +26,18 @@ import pekko.util.{ ByteString, ByteStringBuilder } * INTERNAL API */ @InternalApi private[file] final class ZipArchiveFlowStage( - val shape: FlowShape[ByteString, ByteString]) extends GraphStageLogic(shape) { + val shape: FlowShape[ByteString, ByteString], + deflateCompression: Option[Int] = None) extends GraphStageLogic(shape) { private def in = shape.in private def out = shape.out private val builder = new ByteStringBuilder() private val zip = new ZipOutputStream(builder.asOutputStream) + + override def preStart(): Unit = + deflateCompression.foreach(l => zip.setLevel(l)) + private var emptyStream = true setHandler( @@ -88,7 +93,8 @@ import pekko.util.{ ByteString, ByteStringBuilder } /** * INTERNAL API */ -@InternalApi private[file] final class ZipArchiveFlow extends GraphStage[FlowShape[ByteString, ByteString]] { +@InternalApi private[file] final class ZipArchiveFlow(deflateCompression: Option[Int] = None) + extends GraphStage[FlowShape[ByteString, ByteString]] { val in: Inlet[ByteString] = Inlet(Logging.simpleName(this) + ".in") val out: Outlet[ByteString] = Outlet(Logging.simpleName(this) + ".out") @@ -99,5 +105,5 @@ import pekko.util.{ ByteString, ByteStringBuilder } override val shape: FlowShape[ByteString, ByteString] = FlowShape(in, out) override def createLogic(inheritedAttributes: Attributes): GraphStageLogic = - new ZipArchiveFlowStage(shape) + new ZipArchiveFlowStage(shape, deflateCompression) } diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveManager.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveManager.scala index 9ecd9f524..d19638dda 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveManager.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveManager.scala @@ -25,8 +25,9 @@ import pekko.util.ByteString */ @InternalApi private[file] object ZipArchiveManager { - def zipFlow(): Flow[(ArchiveMetadata, Source[ByteString, Any]), ByteString, NotUsed] = { - val archiveZipFlow = new ZipArchiveFlow() + def zipFlow( + deflateCompression: Option[Int] = None): Flow[(ArchiveMetadata, Source[ByteString, Any]), ByteString, NotUsed] = { + val archiveZipFlow = new ZipArchiveFlow(deflateCompression) Flow[(ArchiveMetadata, Source[ByteString, Any])] .flatMapConcat { case (metadata, stream) => diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala index 16e9d5b78..499ffcdea 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala @@ -33,11 +33,18 @@ object Archive { /** * Flow for compressing multiple files into one ZIP file. */ - def zip(): Flow[Pair[ArchiveMetadata, Source[ByteString, NotUsed]], ByteString, NotUsed] = + def zip( + deflateCompression: Option[Int]): Flow[Pair[ArchiveMetadata, Source[ByteString, NotUsed]], ByteString, NotUsed] = Flow .create[Pair[ArchiveMetadata, Source[ByteString, NotUsed]]]() .map(func(pair => (pair.first, pair.second.asScala))) - .via(scaladsl.Archive.zip().asJava) + .via(scaladsl.Archive.zip(deflateCompression).asJava) + + /** + * Flow for compressing multiple files into one ZIP file. + */ + def zip(): Flow[Pair[ArchiveMetadata, Source[ByteString, NotUsed]], ByteString, NotUsed] = + zip(None) /** * Flow for reading ZIP files. diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala index 5cce29ffc..9bae3e7c1 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala @@ -28,11 +28,18 @@ import java.nio.charset.{ Charset, StandardCharsets } */ object Archive { + /** + * Flow for compressing multiple files into one ZIP file. + */ + def zip( + deflateCompression: Option[Int]): Flow[(ArchiveMetadata, Source[ByteString, Any]), ByteString, NotUsed] = + ZipArchiveManager.zipFlow(deflateCompression) + /** * Flow for compressing multiple files into one ZIP file. */ def zip(): Flow[(ArchiveMetadata, Source[ByteString, Any]), ByteString, NotUsed] = - ZipArchiveManager.zipFlow() + zip(None) /** * Flow for reading ZIP files. diff --git a/file/src/test/scala/docs/scaladsl/ArchiveSpec.scala b/file/src/test/scala/docs/scaladsl/ArchiveSpec.scala index 3103a0035..558084477 100644 --- a/file/src/test/scala/docs/scaladsl/ArchiveSpec.scala +++ b/file/src/test/scala/docs/scaladsl/ArchiveSpec.scala @@ -15,6 +15,7 @@ package docs.scaladsl import java.io._ import java.nio.file.{ Files, Path, Paths } +import java.util.zip.Deflater import org.apache.pekko import pekko.actor.ActorSystem import pekko.stream.connectors.file.ArchiveMetadata @@ -121,6 +122,19 @@ class ArchiveSpec archiveHelper.unzip(pekkoZipped.futureValue).asScala shouldBe inputFiles } + "archive files with compression flag" in { + val inputFiles = generateInputFiles(5, 100) + val inputStream = filesToStream(inputFiles) + val zipFlow = Archive.zip(Some(Deflater.NO_COMPRESSION)) + + val pekkoZipped: Future[ByteString] = + inputStream + .via(zipFlow) + .runWith(Sink.fold(ByteString.empty)(_ ++ _)) + + archiveHelper.unzip(pekkoZipped.futureValue).asScala shouldBe inputFiles + } + "unarchive files" in { val inputFiles = generateInputFiles(5, 100) val inputStream = filesToStream(inputFiles) diff --git a/file/src/test/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlowTest.scala b/file/src/test/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlowTest.scala index f2b8e5e4b..005f8ceab 100644 --- a/file/src/test/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlowTest.scala +++ b/file/src/test/scala/org/apache/pekko/stream/connectors/file/impl/archive/ZipArchiveFlowTest.scala @@ -20,6 +20,7 @@ import pekko.stream.scaladsl.Keep import pekko.stream.testkit.scaladsl.{ TestSink, TestSource } import pekko.testkit.TestKit import pekko.util.ByteString +import java.util.zip.Deflater import org.scalatest.BeforeAndAfterAll import org.scalatest.wordspec.AnyWordSpecLike @@ -50,6 +51,27 @@ class ZipArchiveFlowTest downstream.expectComplete() } } + + "compression flag given and stream ends" should { + "emit element only when downstream requests" in { + val (upstream, downstream) = + TestSource[ByteString]() + .via(new ZipArchiveFlow(Some(Deflater.NO_COMPRESSION))) + .toMat(TestSink())(Keep.both) + .run() + + upstream.sendNext(FileByteStringSeparators.createStartingByteString("test")) + upstream.sendNext(ByteString(1)) + upstream.sendNext(FileByteStringSeparators.createEndingByteString()) + upstream.sendComplete() + + downstream.request(2) + downstream.expectNextN(2) + downstream.request(1) + downstream.expectNextN(1) + downstream.expectComplete() + } + } } override def afterAll(): Unit = { From 3cfe995e13135e875717580ea333634847bd39fd Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Sun, 9 Aug 2026 18:57:50 +0100 Subject: [PATCH 2/2] scaladoc --- .../apache/pekko/stream/connectors/file/javadsl/Archive.scala | 3 +++ .../apache/pekko/stream/connectors/file/scaladsl/Archive.scala | 3 +++ 2 files changed, 6 insertions(+) diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala index 499ffcdea..386cebe25 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/javadsl/Archive.scala @@ -32,6 +32,9 @@ object Archive { /** * Flow for compressing multiple files into one ZIP file. + * @param deflateCompression optional compression level, 0-9, where 0 is no compression and 9 is maximum compression. + * If not specified, the default compression level of the underlying library will be used. + * @since 2.0.0 */ def zip( deflateCompression: Option[Int]): Flow[Pair[ArchiveMetadata, Source[ByteString, NotUsed]], ByteString, NotUsed] = diff --git a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala index 9bae3e7c1..01def3ca2 100644 --- a/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala +++ b/file/src/main/scala/org/apache/pekko/stream/connectors/file/scaladsl/Archive.scala @@ -30,6 +30,9 @@ object Archive { /** * Flow for compressing multiple files into one ZIP file. + * @param deflateCompression optional compression level, 0-9, where 0 is no compression and 9 is maximum compression. + * If not specified, the default compression level of the underlying library will be used. + * @since 2.0.0 */ def zip( deflateCompression: Option[Int]): Flow[(ArchiveMetadata, Source[ByteString, Any]), ByteString, NotUsed] =