Guest User

Untitled

a guest
Sep 10th, 2021
159
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Java 1.93 KB | None | 0 0
  1. public class MainTest {
  2.     private TestStream<SalesOrder> createMainStream() {
  3.         var streamB = TestStream.create(AvroCoder.of(SalesOrder.class))
  4.              .advanceWatermarkTo(Instant.now());
  5.         var counter = 0;
  6.         for (SalesOrder salesOrder : this.jsonMap.subList(0, 10)) {
  7.             counter++;
  8.             streamB = streamB.addElements(TimestampedValue.of(salesOrder, Instant.now()));
  9.         }
  10.         //streamB = streamB.advanceProcessingTime(Duration.standardSeconds(5));
  11.         return streamB.advanceWatermarkToInfinity();
  12.     }
  13.  
  14.     @Test
  15.     public void startMockExecution() {
  16.         final TestStream<SalesOrder> stream = createMainStream();
  17.  
  18.         final PCollectionView<List<Outlet>> outletSideInput = this.pipeline
  19.                 // .apply("fetch outlets", new OutletSideInput())
  20.                 .apply(GenerateSequence.from(0).withRate(1, Duration.standardSeconds(20)))
  21.                 .apply(
  22.                         Window.<Long>into(new GlobalWindows())
  23.                                 .triggering(Repeatedly.forever(AfterProcessingTime.pastFirstElementInPane()))
  24.                                 .discardingFiredPanes())
  25.                 .apply("Read OpeningHours form BQ-Table", ParDo.of(new ReadOutletsFromBQ(new BigQueryService())))
  26.                 .apply("print outlets",ParDo.of(new PrintElem<>(elem -> Integer.toString(elem.size()))))
  27.                 .apply("outlets to view", View.asSingleton());
  28.  
  29.  
  30.         val res = this.pipeline
  31.                 .apply(stream)
  32.                 // .apply(Window.into(FixedWindows.of(Duration.standardSeconds(5))))
  33.                 .apply(ParDo.of(new FilterRelevantSalesOrdersAndOrderDetails()))
  34.                 .apply(ParDo.of(new EnrichSalesOrder(null, outletSideInput))
  35.                         .withSideInputs(outletSideInput)
  36.                 )
  37.                 .apply("print enriched", ParDo.of(new PrintElem<>()));
  38.         this.pipeline.run().waitUntilFinish();
  39.     }
  40. }
Advertisement
Add Comment
Please, Sign In to add comment