Skip to content

Commit f6a490a

Browse files
committed
fix unit tests
1 parent b43071c commit f6a490a

1 file changed

Lines changed: 34 additions & 11 deletions

File tree

runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslatorTest.java

Lines changed: 34 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -620,31 +620,54 @@ public void testDiskSizeGbConfig() throws IOException {
620620
}
621621

622622
@Test
623-
public void testDiskProvisioningTranslation() {
624-
DataflowPipelineOptions options =
625-
PipelineOptionsFactory.create().as(DataflowPipelineOptions.class);
623+
public void testDiskProvisioningTranslation() throws IOException {
624+
DataflowPipelineOptions options = buildPipelineOptions();
626625
options.setDiskProvisionedIops(Long.valueOf(7000));
627626
options.setDiskProvisionedThroughputMibps(Long.valueOf(250));
628-
options.setProject("test-project"); // Required for translator
629627

630-
WorkerPool pool = translateWorkerPool(options);
628+
Pipeline p = buildPipeline(options);
629+
p.traverseTopologically(new RecordingPipelineVisitor());
630+
SdkComponents sdkComponents = createSdkComponents(options);
631+
RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(p, sdkComponents, true);
632+
Job job =
633+
DataflowPipelineTranslator.fromOptions(options)
634+
.translate(
635+
p,
636+
pipelineProto,
637+
sdkComponents,
638+
DataflowRunner.fromOptions(options),
639+
Collections.emptyList())
640+
.getJob();
631641

642+
assertEquals(1, job.getEnvironment().getWorkerPools().size());
643+
WorkerPool pool = job.getEnvironment().getWorkerPools().get(0);
632644
assertEquals(Long.valueOf(7000), pool.getDiskProvisionedIops());
633645
assertEquals(Long.valueOf(250), pool.getDiskProvisionedThroughputMibps());
634646
}
635647

636648
@Test
637-
public void testDiskProvisioningTranslationDefaults() {
638-
DataflowPipelineOptions options =
639-
PipelineOptionsFactory.create().as(DataflowPipelineOptions.class);
640-
options.setProject("test-project"); // Required for translator
649+
public void testDiskProvisioningTranslationDefaults() throws IOException {
650+
DataflowPipelineOptions options = buildPipelineOptions();
641651

642-
WorkerPool pool = translateWorkerPool(options);
652+
Pipeline p = buildPipeline(options);
653+
p.traverseTopologically(new RecordingPipelineVisitor());
654+
SdkComponents sdkComponents = createSdkComponents(options);
655+
RunnerApi.Pipeline pipelineProto = PipelineTranslation.toProto(p, sdkComponents, true);
656+
Job job =
657+
DataflowPipelineTranslator.fromOptions(options)
658+
.translate(
659+
p,
660+
pipelineProto,
661+
sdkComponents,
662+
DataflowRunner.fromOptions(options),
663+
Collections.emptyList())
664+
.getJob();
643665

666+
assertEquals(1, job.getEnvironment().getWorkerPools().size());
667+
WorkerPool pool = job.getEnvironment().getWorkerPools().get(0);
644668
assertNull(pool.getDiskProvisionedIops());
645669
assertNull(pool.getDiskProvisionedThroughputMibps());
646670
}
647-
648671
/** A composite transform that returns an output that is unrelated to the input. */
649672
private static class UnrelatedOutputCreator
650673
extends PTransform<PCollection<Integer>, PCollection<Integer>> {

0 commit comments

Comments
 (0)