damccorm commented on code in PR #40200:
URL: https://github.com/apache/beam/pull/40200#discussion_r4062799593


##########
sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java:
##########
@@ -691,7 +743,10 @@ private class WriteUnshardedTempFilesFn extends 
DoFn<UserT, FileResult<Destinati
     @StartBundle
     public void startBundle(StartBundleContext unused) {
       // Reset state in case of reuse. We need to make sure that each bundle 
gets unique writers.
-      writers = Maps.newHashMap();
+      // LinkedHashMap maintains insertion order so the oldest writer can be 
evicted in O(1) FIFO
+      // order when evictWritersWhenFull is enabled.
+      writers = Maps.newLinkedHashMap();

Review Comment:
   What is the rationale behind FIFO eviction? I can imagine a few strategies, 
but LRU might make more sense here at least.



##########
sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java:
##########
@@ -289,7 +299,10 @@ public WriteFiles<UserT, DestinationT, OutputT> 
withNumShards(
     return toBuilder().setNumShardsProvider(numShardsProvider).build();
   }
 
-  /** Set the maximum number of writers created in a bundle before spilling to 
shuffle. */
+  /**
+   * Set the maximum number of writers kept open in a bundle before spilling 
to shuffle (or evicting
+   * the oldest open writer if {@link #withEvictWritersWhenFull()} is enabled).

Review Comment:
   Now that this decision matrix is becoming more complex, it would be nice to 
be more explicit about tradeoffs. e.g. here: `A higher value here can cause 
more memory consumption, but reduces cost of shuffling records.`
   
   and below for `withEvictWritersWhenFull`: `setting this to True avoids the 
cost of shuffling records, but may lead to smaller output files`



##########
sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java:
##########
@@ -302,6 +315,41 @@ public WriteFiles<UserT, DestinationT, OutputT> 
withMaxNumWritersPerBundle(
     return 
toBuilder().setMaxNumWritersPerBundle(maxNumWritersPerBundle).build();
   }
 
+  /**
+   * Returns a new {@link WriteFiles} that evicts the oldest open writer in 
the bundle (FIFO order,
+   * by flushing and closing it) instead of spilling unwritten records to 
shuffle when {@link
+   * #getMaxNumWritersPerBundle()} is reached.
+   *
+   * <p><b>Warning:</b> This option should only be used when the input {@link 
PCollection} elements
+   * within a bundle are already grouped or ordered by writer keys 
(destination/window/pane), such
+   * that consecutive records belong to the same destination. If the input 
{@link PCollection} rows
+   * arrive in random order across more destinations than {@link 
#getMaxNumWritersPerBundle()},
+   * writers will be repeatedly closed and reopened, creating too many small 
files.
+   *
+   * <p>This option only applies to writes {@link 
#withRunnerDeterminedSharding()}.
+   */
+  public WriteFiles<UserT, DestinationT, OutputT> withEvictWritersWhenFull() {
+    return withEvictWritersWhenFull(true);
+  }
+
+  /**
+   * Set this sink to evict the oldest open writer in the bundle (FIFO order, 
by flushing and
+   * closing it) when {@link #getMaxNumWritersPerBundle()} is reached, instead 
of spilling unwritten
+   * records to shuffle.
+   *
+   * <p><b>Warning:</b> This option should only be used when the input {@link 
PCollection} elements
+   * within a bundle are already grouped or ordered by writer keys 
(destination/window/pane), such

Review Comment:
   I don't think you can really guarantee this within the Beam model - 
PCollections are considered unordered, though some runners may retain some 
ordering by coincidence (the only exception is that some accept a 
`RequiresTimeSortedInput` annotation, but that doesn't apply here).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to