From 8ed07054f78aa25f64e303ec5f13ffa6ae7e701c Mon Sep 17 00:00:00 2001 From: Raghu Angadi Date: Thu, 24 Aug 2017 14:33:28 -0700 Subject: [PATCH] Fix min_timestamp used for KafkaIO watermark. --- .../src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java index 7fb4260313c7..dae4c1d4c1b5 100644 --- a/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java +++ b/sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIO.java @@ -82,6 +82,7 @@ import org.apache.beam.sdk.transforms.SerializableFunction; import org.apache.beam.sdk.transforms.SimpleFunction; import org.apache.beam.sdk.transforms.display.DisplayData; +import org.apache.beam.sdk.transforms.windowing.BoundedWindow; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.PBegin; import org.apache.beam.sdk.values.PCollection; @@ -899,7 +900,7 @@ private static class UnboundedKafkaReader extends UnboundedReader