Skip to content

Commit 9f8e00e

Browse files
authored
Switch streaming engine worker harness based on job settings (#35901)
* Switch streaming engine worker harness based on job settings * Restart StreamingWorkerStatusPages whenever harness type is swithched * Adding a StreamingWorkerHarnessFactoryOutput class to hold harness and its dependencies * Using mock in TC instead of java reflection * Removed null check before setting status page * Ran spotlessApply * Updating test to process work * Adding check directpath experiment is set before switching to FanOutStreamingEngineWorkerHarness
1 parent 659cc4d commit 9f8e00e

6 files changed

Lines changed: 587 additions & 145 deletions

File tree

runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -310,6 +310,9 @@ public Integer create(PipelineOptions options) {
310310
class EnableWindmillServiceDirectPathFactory implements DefaultValueFactory<Boolean> {
311311
@Override
312312
public Boolean create(PipelineOptions options) {
313+
if (ExperimentalOptions.hasExperiment(options, "disable_windmill_service_direct_path")) {
314+
return false;
315+
}
313316
return ExperimentalOptions.hasExperiment(options, "enable_windmill_service_direct_path");
314317
}
315318
}

0 commit comments

Comments
 (0)