flink-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Stephan Ewen <se...@apache.org>
Subject Re: Flink Kafka Consumer Behaviour
Date Thu, 06 Oct 2016 15:30:08 GMT
Hi!

There was an issue in the Kafka 0.9 consumer in Flink concerning
checkpoints. It was relevant mostly for lower-throughput topics /
partitions.

It is fixed in the 1.1.3 release. Can you try out the release candidate and
see if that solves your problem?
See here for details on the release candidate:
http://apache-flink-mailing-list-archive.1008284.n3.nabble.com/VOTE-Release-Apache-Flink-1-1-3-RC1-td13860.html

To test this, set the dependency for the flink-connector-kafka-09 to
"1.1.3" and add the staging repository described in the above link to your
pom.xml.

Thanks,
Stephan


On Tue, Oct 4, 2016 at 5:51 AM, ankitcha <ankitchaudhary123@gmail.com>
wrote:

> Hi Prabhu, cc Stephan, Robert,
>
> I was having similar issues where flink Kafka 09 consumer was not
> committing
> offsets to kafka. After digging into JobManager logs, I found that
> checkpoints were getting expired before getting completed and hence
> "checkpoint completed" message was being ignored.
>
> I increased the checkpoint interval from default 10 mins to 30 mins to
> verify, and then checkpoints were getting finished way before timeout (~12
> mins), and then consumer was correctly updating offsets in kafka.
>
> This seems to be working for us at the moment, and also note this scenarios
> normally happens at the start of the job and the consumer group already has
> some decent lag.
>
> So, you might wanna try increasing checkpoint timeouts and check if that
> solves the issue. You should look for following in the jobmanager logs
>
> [Checkpoint Timer] org.apache.flink.runtime.check
> point.CheckpointCoordinator
> - Checkpoint 37 expired before completing.
> [Checkpoint Timer] org.apache.flink.runtime.check
> point.CheckpointCoordinator
> - Triggering checkpoint 38 @ 1474313373634
> [Checkpoint Timer] org.apache.flink.runtime.check
> point.CheckpointCoordinator
> - Checkpoint 38 expired before completing.
> [Checkpoint Timer] org.apache.flink.runtime.check
> point.CheckpointCoordinator
> - Triggering checkpoint 39 @ 1474313973640
>
> --
> Ankit
>
>
>
> --
> View this message in context: http://apache-flink-user-maili
> ng-list-archive.2336050.n4.nabble.com/Flink-Kafka-Consume
> r-Behaviour-tp8257p9300.html
> Sent from the Apache Flink User Mailing List archive. mailing list archive
> at Nabble.com.
>

Mime
View raw message