This is an automated email from the ASF dual-hosted git repository.
damccorm pushed a commit to branch release-2.77
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/release-2.77 by this push:
new be2778396fb Exposing evictWritersWhenFull in AvroIO (#40336) (#40338)
be2778396fb is described below
commit be2778396fbc0a175de18f4fdd13ec6097566d7c
Author: darshan-sj <[email protected]>
AuthorDate: Wed Sep 30 10:52:03 2026 +0000
Exposing evictWritersWhenFull in AvroIO (#40336) (#40338)
---
.../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 6b23695c21a..3df2feb24fc 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
@@ -615,6 +615,7 @@ public class AvroIO {
.setMetadata(ImmutableMap.of())
.setWindowedWrites(false)
.setNoSpilling(false)
+ .setEvictWritersWhenFull(false)
.setSyncInterval(DataFileConstants.DEFAULT_SYNC_INTERVAL);
}
@@ -1426,6 +1427,8 @@ public class AvroIO {
abstract boolean getNoSpilling();
+ abstract boolean getEvictWritersWhenFull();
+
abstract @Nullable Integer getMaxNumWritersPerBundle();
abstract @Nullable FilenamePolicy getFilenamePolicy();
@@ -1488,6 +1491,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);
@@ -1698,6 +1704,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) {
@@ -1819,6 +1830,9 @@ public class AvroIO {
if (getBadRecordErrorHandler() != null) {
write = write.withBadRecordErrorHandler(getBadRecordErrorHandler());
}
+ if (getEvictWritersWhenFull()) {
+ write = write.withEvictWritersWhenFull();
+ }
return input.apply("Write", write);
}