Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- public class MainTest {
- private TestStream<SalesOrder> createMainStream() {
- var streamB = TestStream.create(AvroCoder.of(SalesOrder.class))
- .advanceWatermarkTo(Instant.now());
- var counter = 0;
- for (SalesOrder salesOrder : this.jsonMap.subList(0, 10)) {
- counter++;
- streamB = streamB.addElements(TimestampedValue.of(salesOrder, Instant.now()));
- }
- //streamB = streamB.advanceProcessingTime(Duration.standardSeconds(5));
- return streamB.advanceWatermarkToInfinity();
- }
- @Test
- public void startMockExecution() {
- final TestStream<SalesOrder> stream = createMainStream();
- final PCollectionView<List<Outlet>> outletSideInput = this.pipeline
- // .apply("fetch outlets", new OutletSideInput())
- .apply(GenerateSequence.from(0).withRate(1, Duration.standardSeconds(20)))
- .apply(
- Window.<Long>into(new GlobalWindows())
- .triggering(Repeatedly.forever(AfterProcessingTime.pastFirstElementInPane()))
- .discardingFiredPanes())
- .apply("Read OpeningHours form BQ-Table", ParDo.of(new ReadOutletsFromBQ(new BigQueryService())))
- .apply("print outlets",ParDo.of(new PrintElem<>(elem -> Integer.toString(elem.size()))))
- .apply("outlets to view", View.asSingleton());
- val res = this.pipeline
- .apply(stream)
- // .apply(Window.into(FixedWindows.of(Duration.standardSeconds(5))))
- .apply(ParDo.of(new FilterRelevantSalesOrdersAndOrderDetails()))
- .apply(ParDo.of(new EnrichSalesOrder(null, outletSideInput))
- .withSideInputs(outletSideInput)
- )
- .apply("print enriched", ParDo.of(new PrintElem<>()));
- this.pipeline.run().waitUntilFinish();
- }
- }
Advertisement
Add Comment
Please, Sign In to add comment