【问题标题】:Kinesis consumer returning empty record (boto, python)Kinesis 消费者返回空记录(boto、python)
【发布时间】:2018-05-09 19:36:56
【问题描述】:

我在检查写入 Kinesis 的数据时遇到问题。看起来下面的示例应该可以工作,但是我得到了一个从 get_records 返回的空列表(在 Records 字段中)。有什么想法会发生什么吗?

import uuid
import boto3
import time


streamname = 'mytestStream'
session = boto3.session.Session() 
kinesis_client = session.client('kinesis', region_name='us-east-1')


##### WRITE TO KINESIS

partitionkey = str(uuid.uuid4())[:8]
put_response = kinesis_client.put_record(StreamName=streamname,Data='mytestdata',PartitionKey=partitionkey)

time.sleep(5)


##### READ FROM KINESIS

shard_id = kinesis_client.describe_stream(StreamName=streamname)['StreamDescription']['Shards'][0]['ShardId']
shard_iterator = kinesis_client.get_shard_iterator(StreamName=streamname, ShardId=shard_id, ShardIteratorType="LATEST")["ShardIterator"]
data_from_kinesis = kinesis_client.get_records(ShardIterator=shard_iterator)

谢谢!

【问题讨论】:

    标签: python amazon-web-services boto3 amazon-kinesis


    【解决方案1】:

    如果您将使用 LATEST 检查点,您应该首先开始读取流,然后放置记录。在您的示例中,时间线如下;

    • 在 t0:流中的最新检查点位于 101。
    • 在 t1(主线程):您将记录放入流中,记录位于检查点 102。
    • 在 t2(主线程):您在 LATEST 点(即 103)开始跟踪流。

    要解决此问题,您应该在不同的线程中运行生产者和消费者。正确的流程应该是这样的;

    • at t0(消费者线程):在 LATEST 位置开始拖尾蒸汽,即 201。
    • 在 t1(生产者线程):您将记录放入流中,并将记录放置在检查点 202 上。
    • 在 t2(消费者线程):随着服务器端的分片向前移动(因为您刚刚添加了数据)并且您从检查点 201 开始一直在分片中,您迭代新的检查点 202 并显示您的数据。李>

    【讨论】:

      猜你喜欢
      • 2015-10-23
      • 1970-01-01
      • 1970-01-01
      • 2017-09-22
      • 2019-05-11
      • 2020-10-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多