使用Boto3(Python)的Kinesis消费者返回空记录问题排查
调用Kinesis get_records返回空列表的排查思路
我之前也碰到过类似的问题,结合你的代码来看,最可能的原因和解决办法如下:
最常见的原因:ShardIteratorType 参数错误
你代码里用了 ShardIteratorType="LATEST",这个类型的迭代器会从当前时刻之后写入的新记录开始读取——也就是说你刚才调用put_record写入的那条数据,已经在这个迭代器的"时间线之前"了,自然读不到。
解决办法很简单,根据你的需求换合适的迭代器类型:
- 如果你想读取流中所有可用的旧记录(包括刚写入的这条),用
TRIM_HORIZON - 如果你只想读取写入那条记录之后的内容,可以用
AFTER_SEQUENCE_NUMBER,并传入put_response['SequenceNumber']作为起始序列号 - 如果需要读取某个时间点之后的记录,用
AT_TIMESTAMP
比如修改你的读取代码:
##### READ FROM KINESIS shard_id = kinesis_client.describe_stream(StreamName=streamname)['StreamDescription']['Shards'][0]['ShardId'] # 改用TRIM_HORIZON读取流中最早的可用记录 shard_iterator = kinesis_client.get_shard_iterator( StreamName=streamname, ShardId=shard_id, ShardIteratorType="TRIM_HORIZON" )["ShardIterator"] data_from_kinesis = kinesis_client.get_records(ShardIterator=shard_iterator)
其他可能的原因
- 分片同步延迟:虽然你sleep了5秒,但极端情况下Kinesis的分片数据同步可能需要更久。可以尝试多调用几次
get_records(每次用返回的NextShardIterator),或者延长等待时间。 - 选错了分片:如果你的流有多个分片,写入的记录可能不在你获取的第一个分片里。因为PartitionKey的哈希值会决定记录进入哪个分片,你可以通过
put_response['ShardId']确认写入的分片ID,再用这个ID去获取迭代器。 - 权限问题:虽然你能调用API,但如果IAM权限没有明确允许读取目标分片的记录,也可能返回空。检查你的IAM策略是否包含
kinesis:GetRecords、kinesis:GetShardIterator和kinesis:DescribeStream权限。
内容的提问来源于stack exchange,提问作者TH22
相关产品推荐
相关产品推荐

