From b7a7dad30c15c35ed49eb34983529395c28998b0 Mon Sep 17 00:00:00 2001 From: Ale Tognola Date: Mon, 29 Jun 2026 19:13:57 +0200 Subject: [PATCH] [IcebergIO] Fix stale openWriters count in RecordWriterManager MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit RecordWriterManager uses a Guava Cache with expireAfterAccess(1 min) for writer eviction. However, Guava Cache eviction is lazy — expired entries are only removed on subsequent access or explicit cleanUp(). The write() method checks `openWriters >= maxNumWriters` without first calling cleanUp(), so expired writers remain in the cache and openWriters stays stale. With DEFAULT_MAX_WRITERS_PER_BUNDLE = 20 and more than 20 partitions, this causes 100% false spill — every record for the 21st+ partition is rejected as if the writer pool is full. Add writers.cleanUp() before the capacity check to evict expired entries and update openWriters accurately. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../org/apache/beam/sdk/io/iceberg/RecordWriterManager.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java index 1e25e6b9c234..d5439ad83831 100644 --- a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java +++ b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/RecordWriterManager.java @@ -163,7 +163,10 @@ boolean write(Record record) { @Nullable RecordWriter writer = writers.getIfPresent(routingPartitionKey); if (writer == null && openWriters >= maxNumWriters) { - return false; + writers.cleanUp(); + if (openWriters >= maxNumWriters) { + return false; + } } writer = fetchWriterForPartition(routingPartitionKey, writer); writer.write(record);