kennknowles commented on code in PR #40321:
URL: https://github.com/apache/beam/pull/40321#discussion_r4134749978
##########
runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java:
##########
@@ -30,6 +34,20 @@
/** Utility class for staging files to GCS. */
public class GcsStager implements Stager {
+ @AutoValue
+ abstract static class StagedFilesCacheKey {
+ abstract String getStagingLocation();
+
+ abstract List<StagedFile> getFilesToStage();
+
+ static StagedFilesCacheKey of(String stagingLocation, List<StagedFile>
filesToStage) {
+ return new AutoValue_GcsStager_StagedFilesCacheKey(stagingLocation,
filesToStage);
+ }
+ }
+
+ private static final ConcurrentHashMap<StagedFilesCacheKey,
List<DataflowPackage>>
+ STAGED_FILES_CACHE = new ConcurrentHashMap<>();
Review Comment:
Maybe use a `Cache` with a size bound that will be big enough for any
reasonable situation, but will prevent future weirdness if things change.
##########
runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java:
##########
@@ -49,6 +67,16 @@ public static GcsStager fromOptions(PipelineOptions options)
{
*/
@Override
public List<DataflowPackage> stageFiles(List<StagedFile> filesToStage) {
+ String stagingLocation = options.getStagingLocation();
+ if (stagingLocation != null &&
TestDataflowRunner.class.equals(options.getRunner())) {
Review Comment:
Why depend on the runner here? Is there a reason it cannot just be always on?
--
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]