Flink消费AWS Kinesis数据源时的线程模型相关理解是否正确?
认知核实结果
- 第一点认知完全正确
Flink Kinesis消费者的shard分配遵循均匀分配规则:当消费者并行度与你当前Kinesis流的shard总数(10个)相等时,每个消费者并行实例(对应1个slot)会恰好分配1个shard,一对一处理对应shard的数据;当消费者并行度小于shard总数时,多出来的shard会被均匀分配给各个并行实例,部分slot会承接多个shard的处理任务。 - 第二点认知完全正确
使用默认的轮询拉取模式时,每个运行Kinesis消费者子任务的slot都会启动一个独立的分片发现线程,持续检测流的shard变更(比如分片分裂、合并产生的新shard)。每个分配到该子任务的shard都会对应独立的消费线程,哪怕单个slot被分配了多个shard,也会为每个shard单独创建消费线程拉取数据,不会出现多个shard复用同一个消费线程的情况。
内容的提问来源于stack exchange,提问作者N A
相关产品推荐
相关产品推荐

