From 686be0bdb56535985847a03b1c8b6347df89ed5a Mon Sep 17 00:00:00 2001 From: mynameborat Date: Mon, 23 Sep 2019 15:39:24 -0700 Subject: [PATCH] SAMZA-2329: Chain watermark future to the result future in onMessageAsync --- .../samza/operators/impl/OperatorImpl.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImpl.java b/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImpl.java index 3d32be3197..528acc66f2 100644 --- a/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImpl.java +++ b/samza-core/src/main/java/org/apache/samza/operators/impl/OperatorImpl.java @@ -191,14 +191,14 @@ public final CompletionStage onMessageAsync(M message, MessageCollector co .toArray(CompletableFuture[]::new)); }); - return result.thenAccept(x -> { - WatermarkFunction watermarkFn = getOperatorSpec().getWatermarkFn(); - if (watermarkFn != null) { - // check whether there is new watermark emitted from the user function - Long outputWm = watermarkFn.getOutputWatermark(); - propagateWatermark(outputWm, collector, coordinator); - } - }); + WatermarkFunction watermarkFn = getOperatorSpec().getWatermarkFn(); + if (watermarkFn != null) { + // check whether there is new watermark emitted from the user function + Long outputWm = watermarkFn.getOutputWatermark(); + return result.thenCompose(ignored -> propagateWatermark(outputWm, collector, coordinator)); + } + + return result; } /**