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);