bharadwaj-aditya commented on code in PR #40200:
URL: https://github.com/apache/beam/pull/40200#discussion_r4062676661
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java:
##########
@@ -752,8 +799,63 @@ public void processElement(
writeOrClose(writer, formattedRecord);
}
+ private Writer<DestinationT, OutputT> openAndRegisterWriter(
+ WriterKey<DestinationT> key,
+ BoundedWindow window,
+ PaneInfo paneInfo,
+ DestinationT destination)
+ throws Exception {
+ String uuid = UUID.randomUUID().toString();
+ LOG.info(
+ "Opening writer {} for window {} pane {} destination {}",
+ uuid,
+ window,
+ paneInfo,
+ destination);
+ Writer<DestinationT, OutputT> writer = writeOperation.createWriter();
+ writer.setDestination(destination);
+ writer.open(uuid);
+ writers.put(key, writer);
+ LOG.debug("Done opening writer");
+ return writer;
+ }
+
+ private void evictOldestWriter() throws Exception {
+ Iterator<Map.Entry<WriterKey<DestinationT>, Writer<DestinationT,
OutputT>>> iterator =
+ writers.entrySet().iterator();
+ Map.Entry<WriterKey<DestinationT>, Writer<DestinationT, OutputT>>
eldestEntry =
+ iterator.next();
+ iterator.remove();
Review Comment:
This method needs to be thread safe with a check on number of writers. If
not it is possible that the list is empty if multiple threads enter this block
together
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/io/WriteFiles.java:
##########
@@ -752,8 +799,63 @@ public void processElement(
writeOrClose(writer, formattedRecord);
}
+ private Writer<DestinationT, OutputT> openAndRegisterWriter(
Review Comment:
this needs to be thread safe.
--
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]