Eugene Kirpichov reassigned BEAM-3696:

    Assignee: Jean-Baptiste Onofré  (was: Reuven Lax)

> MQTT IO should compute watermark and ack messages outside of 
> finalizeCheckpoint method
> --------------------------------------------------------------------------------------
>                 Key: BEAM-3696
>                 URL: https://issues.apache.org/jira/browse/BEAM-3696
>             Project: Beam
>          Issue Type: Bug
>          Components: sdk-java-extensions
>    Affects Versions: 2.2.0
>         Environment: - Flink - beam-runners-flink_2.10:2.2.0
> - Beam and related jars - 2.2.0
>            Reporter: Maxim Kolchin
>            Assignee: Jean-Baptiste Onofré
>            Priority: Major
> I'm experiencing a situation when an incoming message isn't acknowledged 
> (therefore in sometime broker resend it) and the watermark is not updated 
> while new messages are coming continuously.
> After some time I've discovered that this situation is related to the fact 
> that finalizaCheckpoint is not being called.
> I took a look at the Pubsub IO implementation and found that they expect such 
> situation and do not compute watermark and ack messages in 
> finalizeCheckpoint. Here is the comment about that: 
> [https://github.com/apache/beam/blob/master/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/pubsub/PubsubUnboundedSource.java#L289]
> Should MQTT IO do the same?

This message was sent by Atlassian JIRA

Reply via email to