添加新分片后KCL Lease表未更新,导致Worker闲置问题求助
问题
我正在使用版本为2.1.3的amazon-kclpy(基于MultiLangDaemon的Python库)。根据AWS文档,当为流添加新分片或启动新Worker(我这边是Kubernetes中的Pod)时,应触发分片重分配事件,之后Lease表中应出现新记录。但实际在这两种场景下,Lease表都未生成新记录,导致新增的Worker处于闲置状态并进入睡眠循环,相关日志如下:
2023-12-24 17:46:34,750 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - Initializing LeaseCoordinator attempt 1 2023-12-24 17:46:37,860 [multi-lang-daemon-0000] INFO s.a.k.leases.LeaseCleanupManager [NONE] - Starting lease cleanup thread. 2023-12-24 17:46:37,861 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - Starting LeaseCoordinator 2023-12-24 17:46:37,861 [pool-15-thread-1] INFO s.a.k.leases.LeaseCleanupManager [NONE] - Number of pending leases to clean before the scan : 0 2023-12-24 17:46:37,926 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - Scheduling periodicShardSync 2023-12-24 17:46:37,928 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - Initialization complete. Starting worker loop. 2023-12-24 17:46:38,026 [multi-lang-daemon-0000] INFO s.a.k.c.DeterministicShuffleShardSyncLeaderDecider [NONE] - Elected leaders: fe86945c-96d7-423e-a157-97c5c6ad6f79 2023-12-24 17:47:05,036 [multi-lang-daemon-0000] INFO s.a.k.c.DiagnosticEventLogger [NONE] - Current thread pool executor state: ExecutorStateEvent(executorName=SchedulerThreadPoolExecutor, currentQueueSize=0, activeThreads=0, coreThreads=0, leasesOwned=0, largestPoolSize=0, maximumPoolSize=2147483647) 2023-12-24 17:47:35,044 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - No activities assigned 2023-12-24 17:47:35,045 [multi-lang-daemon-0000] INFO s.a.k.c.DiagnosticEventLogger [NONE] - Current thread pool executor state: ExecutorStateEvent(executorName=SchedulerThreadPoolExecutor, currentQueueSize=0, activeThreads=0, coreThreads=0, leasesOwned=0, largestPoolSize=0, maximumPoolSize=2147483647) 2023-12-24 17:47:35,045 [multi-lang-daemon-0000] INFO s.a.kinesis.coordinator.Scheduler [NONE] - Sleeping ...
请问这是权限问题还是Python代码存在问题?
排查方向
1. 权限问题
- 检查Pod绑定的IAM角色是否包含以下必要权限:
kinesis:DescribeStream/kinesis:DescribeStreamSummary:Worker需要获取Kinesis流的分片元数据dynamodb:PutItem、dynamodb:GetItem、dynamodb:UpdateItem、dynamodb:Scan:对Lease表的读写及扫描权限- 确认权限策略中的资源ARN是否正确,比如Lease表的ARN是否完整匹配,没有限制特定条目
2. 分片同步配置问题
- 检查
periodicShardSyncIntervalMillis参数:默认5分钟,若配置过大,新分片/Worker的同步会延迟 - 确保所有Worker使用同一个Lease表:不同表会导致分片分配逻辑隔离,新Worker无法感知已有分片
- 验证Leader节点状态:日志显示已选举Leader,但需确认该Leader节点能正常访问Kinesis流和Lease表,未被资源限制或网络隔离
3. Amazon-KCLPy配置校验
- 确认所有Worker的
applicationName配置一致:这个参数是Lease表关联的核心标识,不一致会导致无法共享分片信息 - 检查
shardSyncEnabled是否启用:默认开启,若被禁用会直接阻止分片同步逻辑执行 - 核对
kinesisStreamName/kinesisStreamArn配置:确保流名称/ARN拼写正确,没有环境差异导致的配置错误
4. 代码逻辑问题
- 检查自定义RecordProcessor的
initialize、shardEnded方法:若存在阻塞、未正确处理异常等情况,会导致Lease更新失败,影响分片分配 - 确认未手动修改Lease表记录:手动删除或修改分片对应的Lease条目会破坏KCL的自动分配机制,导致新分片无法被识别
内容的提问来源于stack exchange,提问作者omer bar lev
相关产品推荐
相关产品推荐

