This is an automated email from the ASF dual-hosted git repository.
pjfanning pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/pekko-connectors.git
The following commit(s) were added to refs/heads/main by this push:
new 817f614e3 S3: release disk-buffered chunks when the upload is done
with them (#1930)
817f614e3 is described below
commit 817f614e39f0d9e2beb753001bbc0c9a6bb9cc34
Author: PJ Fanning <[email protected]>
AuthorDate: Wed Sep 30 14:12:30 2026 +0100
S3: release disk-buffered chunks when the upload is done with them (#1930)
* S3: release disk-buffered chunks when the upload is done with them
DiskBuffer deleted its temp file only once the chunk had been materialized
as
many times as the retry budget allows, so a chunk that uploaded first time -
two materializations out of eight - was never deleted, and survived on disk
until the JVM exited via deleteOnExit.
Give Chunk a dispose() that DiskChunk implements as a file delete, and call
it
from the RetryFlow decision branch that declines to retry, which is the
point
the chunk is definitively finished with. Also delete the temp file in
postStop
when the stage never emitted a chunk. The existing countdown stays as the
net
for retries being exhausted, which is precisely when it fires.
Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
* S3: delete disk-buffered temp files in code, not at JVM exit
deleteOnExit only fires on an orderly shutdown, never on SIGKILL, and its
JVM-global registration is never pruned. Replace it with a registry, held by
the DiskBuffer for the upload it belongs to, that tracks the temp file of
every
chunk it emits.
A chunk releases its own file as soon as the upload is finished with it. Any
file still registered when the upload stream terminates - a chunk emitted
and
then abandoned by a cancelled or failed upload - is deleted by a cleanUp()
hung
off watchTermination, so nothing waits for the JVM to exit.
Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
---------
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
.../pekko/stream/connectors/s3/impl/Chunk.scala | 28 +++++++++++-
.../stream/connectors/s3/impl/DiskBuffer.scala | 50 +++++++++++++++++++---
.../stream/connectors/s3/impl/MemoryBuffer.scala | 2 +-
.../pekko/stream/connectors/s3/impl/S3Stream.scala | 18 ++++++--
.../stream/connectors/s3/impl/DiskBufferSpec.scala | 48 +++++++++++++++++++++
5 files changed, 135 insertions(+), 11 deletions(-)
diff --git
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
index 237a432b6..5648447df 100644
--- a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
+++ b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/Chunk.scala
@@ -18,18 +18,44 @@ import pekko.stream.scaladsl.Source
import pekko.NotUsed
import pekko.annotation.InternalApi
import pekko.http.scaladsl.model.{ ContentTypes, HttpEntity, RequestEntity }
+import pekko.stream.FlowShape
+import pekko.stream.stage.GraphStage
import pekko.util.ByteString
+import java.io.File
+
/**
* Internal Api
*/
@InternalApi private[impl] sealed trait Chunk {
def asEntity(): RequestEntity
def size: Int
+
+ /**
+ * Releases whatever backs this chunk, once it is known that `data` will not
be materialized again.
+ * Must not be called while an upload attempt for the chunk may still be
retried.
+ */
+ def dispose(): Unit = ()
}
-@InternalApi private[impl] final case class DiskChunk(data: Source[ByteString,
NotUsed], size: Int) extends Chunk {
+@InternalApi private[impl] final class DiskChunk(val data: Source[ByteString,
NotUsed],
+ val size: Int,
+ file: File,
+ registry: TempFileRegistry) extends Chunk {
def asEntity(): RequestEntity =
HttpEntity(ContentTypes.`application/octet-stream`, size, data)
+ override def dispose(): Unit = registry.release(file)
+}
+
+/**
+ * A stage that buffers a chunk, and can release whatever any chunk it emitted
still holds.
+ */
+@InternalApi private[impl] trait ChunkBuffer extends
GraphStage[FlowShape[ByteString, Chunk]] {
+
+ /**
+ * Releases the resources of every chunk this buffer emitted that was not
disposed of individually.
+ * Called once the stream the chunks were emitted into has terminated,
however it terminated.
+ */
+ def cleanUp(): Unit = ()
}
@InternalApi private[impl] final case class MemoryChunk(data: ByteString)
extends Chunk {
diff --git
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
index 75fc824a6..14841584c 100644
---
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
+++
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBuffer.scala
@@ -15,8 +15,10 @@ package org.apache.pekko.stream.connectors.s3.impl
import java.io.{ File, FileOutputStream }
import java.nio.BufferOverflowException
+import java.io.File
import java.nio.file.Files
import java.nio.file.Path
+import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicInteger
import org.apache.pekko
@@ -28,7 +30,6 @@ import pekko.stream.FlowShape
import pekko.stream.Inlet
import pekko.stream.Outlet
import pekko.stream.scaladsl.FileIO
-import pekko.stream.stage.GraphStage
import pekko.stream.stage.GraphStageLogic
import pekko.stream.stage.InHandler
import pekko.stream.stage.OutHandler
@@ -36,6 +37,29 @@ import pekko.util.ByteString
import scala.concurrent.ExecutionContext
+/**
+ * Internal Api
+ *
+ * Tracks the temp files of chunks that have been emitted but not yet
released, so that they can be deleted
+ * as soon as they are no longer needed, and at the latest when the stream
they were emitted into ends.
+ */
+@InternalApi private[impl] final class TempFileRegistry {
+ private val files = ConcurrentHashMap.newKeySet[File]()
+
+ def register(file: File): Unit = { files.add(file): Unit }
+
+ /** Deletes the file and stops tracking it. Deleting a file that is already
gone is a no-op. */
+ def release(file: File): Unit = {
+ files.remove(file)
+ file.delete(): Unit
+ }
+
+ def releaseAll(): Unit = {
+ files.forEach { file => file.delete(): Unit }
+ files.clear()
+ }
+}
+
/**
* Internal Api
*
@@ -44,14 +68,20 @@ import scala.concurrent.ExecutionContext
* The stage waits for the incoming stream to complete. After that, it emits a
single Chunk item on its output. The Chunk
* contains a bytestream source that can be materialized multiple times, and
the total size of the file.
*
- * @param maxMaterializations Number of expected materializations for the
completed chunk. After this, the temp file is deleted.
+ * @param maxMaterializations Maximum number of materializations the completed
chunk may see, which is reached only
+ * when every upload retry is used. After this, the
temp file is deleted. In the ordinary
+ * case the chunk is disposed of as soon as the
upload it belongs to is finished with it.
* @param maxSize Maximum size on disk to buffer
*/
@InternalApi private[impl] final class DiskBuffer(maxMaterializations: Int,
maxSize: Int, tempPath: Option[Path])
- extends GraphStage[FlowShape[ByteString, Chunk]] {
+ extends ChunkBuffer {
require(maxMaterializations > 0, "maxMaterializations should be at least 1")
require(maxSize > 0, "maximumSize should be at least 1")
+ private val registry = new TempFileRegistry
+
+ override def cleanUp(): Unit = registry.releaseAll()
+
val in = Inlet[ByteString]("DiskBuffer.in")
val out = Outlet[Chunk]("DiskBuffer.out")
override val shape = FlowShape.of(in, out)
@@ -65,7 +95,6 @@ import scala.concurrent.ExecutionContext
.map(dir => Files.createTempFile(dir, "s3-buffer-", ".bin"))
.getOrElse(Files.createTempFile("s3-buffer-", ".bin"))
.toFile
- path.deleteOnExit()
var length = 0
val pathOut = new FileOutputStream(path)
@@ -82,6 +111,8 @@ import scala.concurrent.ExecutionContext
pull(in)
}
+ private var emitted = false
+
override def onUpstreamFinish(): Unit = {
if (isAvailable(out)) emit()
completeStage()
@@ -92,20 +123,27 @@ import scala.concurrent.ExecutionContext
try {
pathOut.close()
} catch { case x: Throwable => () }
+ finally {
+ // nothing downstream can reference the file if the chunk was never
emitted
+ if (!emitted) path.delete(): Unit
+ }
private def emit(): Unit = {
pathOut.close()
+ emitted = true
+
+ registry.register(path)
val deleteCounter = new AtomicInteger(maxMaterializations)
val src = FileIO.fromPath(path.toPath, 65536).mapMaterializedValue { f
=>
if (deleteCounter.decrementAndGet() <= 0)
f.onComplete { _ =>
- path.delete()
+ registry.release(path)
}(ExecutionContext.parasitic)
NotUsed
}
- emit(out, DiskChunk(src, length), () => completeStage())
+ emit(out, new DiskChunk(src, length, path, registry), () =>
completeStage())
}
setHandlers(in, out, this)
}
diff --git
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
index 682116f68..3646762f6 100644
---
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
+++
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/MemoryBuffer.scala
@@ -29,7 +29,7 @@ import pekko.util.ByteString
*
* @param maxSize Maximum size to buffer
*/
-@InternalApi private[impl] final class MemoryBuffer(maxSize: Int) extends
GraphStage[FlowShape[ByteString, Chunk]] {
+@InternalApi private[impl] final class MemoryBuffer(maxSize: Int) extends
ChunkBuffer {
val in = Inlet[ByteString]("MemoryBuffer.in")
val out = Outlet[Chunk]("MemoryBuffer.out")
override val shape = FlowShape.of(in, out)
diff --git
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
index ce974da44..9c79c7385 100644
---
a/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
+++
b/s3/src/main/scala/org/apache/pekko/stream/connectors/s3/impl/S3Stream.scala
@@ -1126,7 +1126,8 @@ import scala.util.{ Failure, Success, Try }
initialUploadState: Option[(String, Int)] = None)(
parallelism: Int): Flow[ByteString, UploadPartResponse, NotUsed] = {
- def getChunkBuffer(chunkSize: Int, bufferSize: Int, maxRetriesPerChunk:
Int)(implicit settings: S3Settings) =
+ def getChunkBuffer(chunkSize: Int, bufferSize: Int, maxRetriesPerChunk:
Int)(
+ implicit settings: S3Settings): ChunkBuffer =
settings.bufferType match {
case MemoryBufferType =>
new MemoryBuffer(bufferSize)
@@ -1185,8 +1186,10 @@ import scala.util.{ Failure, Success, Try }
import conf.multipartUploadSettings.retrySettings._
+ val chunkBuffer = getChunkBuffer(chunkSize, chunkBufferSize,
maxRetries) // creates the chunks
+
SplitAfterSize(chunkSize, chunkBufferSize)(atLeastOneByteString)
- .via(getChunkBuffer(chunkSize, chunkBufferSize, maxRetries)) //
creates the chunks
+ .via(chunkBuffer)
.mergeSubstreamsWithParallelism(parallelism)
.filter(_.size > 0)
.via(atLeastOne)
@@ -1198,8 +1201,11 @@ import scala.util.{ Failure, Success, Try }
if (isTransientError(r.status)) {
r.entity.discardBytes()
Some(chunkAndUploadInfo)
- } else
+ } else {
+ // the request has been sent and will not be retried, so the
buffered chunk is no longer needed
+ chunkAndUploadInfo._1.dispose()
None
+ }
case (chunkAndUploadInfo, (Failure(_), _)) =>
// Treat any exception as transient.
Some(chunkAndUploadInfo)
@@ -1209,6 +1215,12 @@ import scala.util.{ Failure, Success, Try }
handleChunkResponse(response, upload, index,
conf.multipartUploadSettings.retrySettings)
}
.mergeSubstreamsWithParallelism(parallelism)
+ .watchTermination((_, done) => {
+ // a chunk that was emitted and then abandoned - a cancelled or
failed upload - is never disposed
+ // of individually, so release whatever is still held once this
upload has finished either way
+ done.onComplete(_ =>
chunkBuffer.cleanUp())(ExecutionContext.parasitic)
+ NotUsed
+ })
}
.mapMaterializedValue(_ => NotUsed)
}
diff --git
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
index 16651bfdb..5f8358708 100644
---
a/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
+++
b/s3/src/test/scala/org/apache/pekko/stream/connectors/s3/impl/DiskBufferSpec.scala
@@ -95,4 +95,52 @@ class DiskBufferSpec(_system: ActorSystem)
}
}
+
+ it should "delete its temp file when the chunk is disposed" in {
+ val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+ val before = tmpDir.list().size
+ val chunk = Source(Vector(ByteString(1, 2, 3)))
+ .via(new DiskBuffer(8, 200, Some(tmpDir.toPath)))
+ .runWith(Sink.seq)
+ .futureValue
+ .head
+
+ // a successful upload materializes the chunk far fewer times than the
retry budget allows,
+ // so disposal, not the materialization count, is what releases the file
+ chunk.asInstanceOf[DiskChunk].data.runWith(Sink.ignore).futureValue
+ tmpDir.list().size should be(before + 1)
+
+ chunk.dispose()
+ tmpDir.list().size should be(before)
+ }
+
+ it should "delete the temp files of chunks that were never disposed of when
it is cleaned up" in {
+ val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+ val before = tmpDir.list().size
+ val buffer = new DiskBuffer(8, 200, Some(tmpDir.toPath))
+
+ Source(Vector(ByteString(1, 2,
3))).via(buffer).runWith(Sink.seq).futureValue
+
+ // the chunk was emitted and then abandoned, as it would be by a cancelled
upload
+ tmpDir.list().size should be(before + 1)
+
+ buffer.cleanUp()
+ tmpDir.list().size should be(before)
+ }
+
+ it should "delete its temp file if it fails before emitting a chunk" in {
+ val tmpDir = Files.createTempDirectory("DiskBufferSpec").toFile()
+ val before = tmpDir.list().size
+
+ Source
+ .failed(new RuntimeException("boom"))
+ .via(new DiskBuffer(8, 200, Some(tmpDir.toPath)))
+ .runWith(Sink.seq)
+ .failed
+ .futureValue shouldBe a[RuntimeException]
+
+ eventually {
+ tmpDir.list().size should be(before)
+ }
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]