flink-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From "Tzu-Li (Gordon) Tai" <tzuli...@apache.org>
Subject Re: Kafka Consumers Partition Discovery doesn't work
Date Thu, 22 Mar 2018 15:09:40 GMT

I think you’ve made a good point: there is currently no logs that tell anything about discovering
a new partition. We should probably add this.

And yes, it would be great if you can report back on this using either the latest master,
release-1.5 or release-1.4 branches.

On 22 March 2018 at 10:24:09 PM, Juho Autio (juho.autio@rovio.com) wrote:

Thanks, that sounds promising. I don't know how to check if it's consuming all partitions?
For example I couldn't find any logs about discovering a new partition. However, did I understand
correctly that this is also fixed in Flink dev? If yes, I could rebuild my 1.5-SNAPSHOT and
try again.

On Thu, Mar 22, 2018 at 4:18 PM, Tzu-Li (Gordon) Tai <tzulitai@apache.org> wrote:
Hi Juho,

Can you confirm that the new partition is consumed, but only that Flink’s reported metrics
do not include them?
If yes, then I think your observations can be explained by this issue: https://issues.apache.org/jira/browse/FLINK-8419

This issue should have been fixed in the recently released 1.4.2 version.


On 22 March 2018 at 8:02:40 PM, Juho Autio (juho.autio@rovio.com) wrote:

According to the docs*, flink.partition-discovery.interval-millis can be set to enable automatic
partition discovery.

I'm testing this, apparently it doesn't work.

I'm using Flink Version: 1.5-SNAPSHOT Commit: 8395508 and FlinkKafkaConsumer010.

I had my flink stream running, consuming an existing topic with 3 partitions, among some other
I modified partitions of an existing topic: 3 -> 4**.
I checked consumer offsets by secor: it's now consuming all 4 partitions.
I checked consumer offset by my flink stream: it's still consuming only the 3 original partitions.

I also checked the Task Metrics of this job from Flink UI and it only offers Kafka related
metrics to be added for 3 partitions (0,1 & 2).

According to Flink UI > Job Manager > Configuration:
– so that's just 1 minute. It's already more than 20 minutes since I added the new partition,
so Flink should've picked it up.

How to debug?

Btw, this job has external checkpoints enabled, done once per minute. Those are also succeeding.

*) https://ci.apache.org/projects/flink/flink-docs-master/dev/connectors/kafka.html#kafka-consumers-topic-and-partition-discovery


~/kafka$ bin/kafka-topics.sh --zookeeper $ZOOKEEPER_HOST --describe --topic my_topic
Topic:my_topic PartitionCount:3 ReplicationFactor:1 Configs:

~/kafka$ bin/kafka-topics.sh --zookeeper $ZOOKEEPER_HOST --alter --topic my_topic --partitions
Adding partitions succeeded!

View raw message