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

damccorm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new c965c6eabeb Exposing evictWritersWhenFull in AvroIO (#40336)
c965c6eabeb is described below

commit c965c6eabeb923c9eaacd860bfc69c4f8d4fc366
Author: darshan-sj <[email protected]>
AuthorDate: Tue Sep 29 16:04:24 2026 +0000

    Exposing evictWritersWhenFull in AvroIO (#40336)
---
 .../core/src/main/java/org/apache/beam/sdk/io/FileIO.java | 15 +++++++++++++++
 .../org/apache/beam/sdk/extensions/avro/io/AvroIO.java    | 14 ++++++++++++++
 2 files changed, 29 insertions(+)

diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileIO.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileIO.java
index e0ff0d51572..3f428900361 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileIO.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/io/FileIO.java
@@ -396,6 +396,7 @@ public class FileIO {
         .setIgnoreWindowing(false)
         .setAutoSharding(false)
         .setNoSpilling(false)
+        .setEvictWritersWhenFull(false)
         .build();
   }
 
@@ -410,6 +411,7 @@ public class FileIO {
         .setIgnoreWindowing(false)
         .setAutoSharding(false)
         .setNoSpilling(false)
+        .setEvictWritersWhenFull(false)
         .build();
   }
 
@@ -1059,6 +1061,8 @@ public class FileIO {
 
     abstract boolean getNoSpilling();
 
+    abstract boolean getEvictWritersWhenFull();
+
     abstract @Nullable Integer getMaxNumWritersPerBundle();
 
     abstract @Nullable Integer getBatchSize();
@@ -1117,6 +1121,8 @@ public class FileIO {
 
       abstract Builder<DestinationT, UserT> setNoSpilling(boolean noSpilling);
 
+      abstract Builder<DestinationT, UserT> setEvictWritersWhenFull(boolean 
evictWritersWhenFull);
+
       abstract Builder<DestinationT, UserT> setMaxNumWritersPerBundle(
           @Nullable Integer maxNumWritersPerBundle);
 
@@ -1366,6 +1372,11 @@ public class FileIO {
       return toBuilder().setNoSpilling(true).build();
     }
 
+    /** See {@link WriteFiles#withEvictWritersWhenFull()}. */
+    public Write<DestinationT, UserT> withEvictWritersWhenFull() {
+      return toBuilder().setEvictWritersWhenFull(true).build();
+    }
+
     /**
      * Set the maximum number of writers created in a bundle before spilling 
to shuffle. See {@link
      * WriteFiles#withMaxNumWritersPerBundle()}.
@@ -1511,6 +1522,7 @@ public class FileIO {
       resolvedSpec.setIgnoreWindowing(getIgnoreWindowing());
       resolvedSpec.setAutoSharding(getAutoSharding());
       resolvedSpec.setNoSpilling(getNoSpilling());
+      resolvedSpec.setEvictWritersWhenFull(getEvictWritersWhenFull());
       if (getMaxNumWritersPerBundle() != null) {
         resolvedSpec.setMaxNumWritersPerBundle(getMaxNumWritersPerBundle());
       }
@@ -1535,6 +1547,9 @@ public class FileIO {
       if (getNoSpilling()) {
         writeFiles = writeFiles.withNoSpilling();
       }
+      if (getEvictWritersWhenFull()) {
+        writeFiles = writeFiles.withEvictWritersWhenFull();
+      }
       if (getMaxNumWritersPerBundle() != null) {
         writeFiles = 
writeFiles.withMaxNumWritersPerBundle(getMaxNumWritersPerBundle());
       }
diff --git 
a/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
 
b/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
index 33b78abad4f..fcb47ba74ab 100644
--- 
a/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
+++ 
b/sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
@@ -612,6 +612,7 @@ public class AvroIO {
         .setMetadata(ImmutableMap.of())
         .setWindowedWrites(false)
         .setNoSpilling(false)
+        .setEvictWritersWhenFull(false)
         .setSyncInterval(DataFileConstants.DEFAULT_SYNC_INTERVAL);
   }
 
@@ -1436,6 +1437,8 @@ public class AvroIO {
 
     abstract boolean getNoSpilling();
 
+    abstract boolean getEvictWritersWhenFull();
+
     abstract @Nullable Integer getMaxNumWritersPerBundle();
 
     abstract @Nullable FilenamePolicy getFilenamePolicy();
@@ -1498,6 +1501,9 @@ public class AvroIO {
 
       abstract Builder<UserT, DestinationT, OutputT> setNoSpilling(boolean 
noSpilling);
 
+      abstract Builder<UserT, DestinationT, OutputT> setEvictWritersWhenFull(
+          boolean evictWritersWhenFull);
+
       abstract Builder<UserT, DestinationT, OutputT> setMaxNumWritersPerBundle(
           @Nullable Integer maxNumWritersPerBundle);
 
@@ -1708,6 +1714,11 @@ public class AvroIO {
       return toBuilder().setNoSpilling(true).build();
     }
 
+    /** See {@link WriteFiles#withEvictWritersWhenFull()}. */
+    public TypedWrite<UserT, DestinationT, OutputT> withEvictWritersWhenFull() 
{
+      return toBuilder().setEvictWritersWhenFull(true).build();
+    }
+
     /** See {@link WriteFiles#withMaxNumWritersPerBundle()}. */
     public TypedWrite<UserT, DestinationT, OutputT> withMaxNumWritersPerBundle(
         @Nullable Integer maxNumWritersPerBundle) {
@@ -1824,6 +1835,9 @@ public class AvroIO {
       if (maxNumWritersPerBundle != null) {
         write = write.withMaxNumWritersPerBundle(maxNumWritersPerBundle);
       }
+      if (getEvictWritersWhenFull()) {
+        write = write.withEvictWritersWhenFull();
+      }
       ErrorHandler<BadRecord, ?> badRecordErrorHandler = 
getBadRecordErrorHandler();
       if (badRecordErrorHandler != null) {
         write = write.withBadRecordErrorHandler(badRecordErrorHandler);

Reply via email to