【发布时间】:2016-10-12 15:14:31
【问题描述】:
我们多次进行了以下测试,但我们很难找到一个合理的解释来解释为什么会发生这种情况:
- 我们创建一个消费者,等待它准备好
- 我们在消费者正在收听的流上发布两条记录
我们有一个分片,有时消费者不会收到有关记录的通知。我们使用不同的 workerId,但具有相同 ApplicationName 的应用程序可能会窃取记录。
KCL 消费者从未获得刚刚发布的记录的原因是什么?
【问题讨论】:
标签: amazon-kinesis
我们多次进行了以下测试,但我们很难找到一个合理的解释来解释为什么会发生这种情况:
我们有一个分片,有时消费者不会收到有关记录的通知。我们使用不同的 workerId,但具有相同 ApplicationName 的应用程序可能会窃取记录。
KCL 消费者从未获得刚刚发布的记录的原因是什么?
【问题讨论】:
标签: amazon-kinesis
我的代码的问题是,一旦您第一次在 KCL 中创建应用程序,Dynamo 中用于检查点的匹配条目会存储分片迭代器的最后位置。
当您重新启动应用程序时,分片迭代器可能会比当前时间滞后很多,并且在您到达发布记录的当前位置之前可能需要进行大量迭代。
因此,出于测试目的,您可能需要每次都创建一个新应用并每晚清理您的发电机表
【讨论】:
您似乎在每个应用程序名称中使用了许多工作人员,但您的流只有一个分片。每个分片只能有一个工作人员处于活动状态并检索数据,其他所有工作人员都将处于空闲状态。
您应确保工作人员的数量不超过分片的数量(或仅增加几个以提高弹性)。
使用 Kinesis,您可以为每个应用程序名称设置多个工作器,以便并行从分片中检索数据。假定具有相同应用程序名称的工作人员在同一个流上一起工作。
但是,每个分片每次只能由一名工作人员访问,例如具有相同应用程序名称的两个工作人员不能同时从同一个分片中读取。如果您有 N 个工作人员但只有 1 个分片,则始终有 1 个工作人员在运行,而 N-1 将处于空闲状态。哪个是活跃的工作人员,是一个调度问题。
如果您希望多个工作人员并行使用来自同一个分片的数据,则需要运行应用程序的另一个实例,但应用程序名称不同。第二个实例被认为是一个完全独立的应用程序,它也在同一个流上运行。
KCL 使用 DynamoDB 处理每个应用程序名称的状态信息(例如,分片数、工作程序、检查点、工作程序分片映射等)。每个区域的每个应用程序名称都应该是唯一的,并且有自己的 DynamoDB 表。
当工作人员启动时,它会查询 DynamoDB 以查看应用程序名称的表是否存在。如果没有表,它会创建一个新表并将其状态写入其中。工作人员自动发现分片并创建处理器来处理来自它们的数据。
worker 可以从许多分片中读取数据,例如它使用 DynamoDB 表来跟踪它检索到的记录,并将使用它来获取新的分片迭代器。但是,正如我之前所说,一个分片每次只能由一个工作人员访问(当然我指的是具有相同应用程序名称的工作人员)。
如果第二个工作器启动,它会查询 DynamoDB,这一次发现应用程序和流的表存在。所以它只是通过在表中写入状态来注册自己。通过 DynamoDB 表,工作人员可以发现彼此并划分工作,例如如果存在多于一个分片,则每个工人都可以从一半分片中读取。否则,新工作人员将处于空闲状态。
最后一件事。如果我没记错的话,将一条记录发送到 kinesis 并准备好供工作人员使用可能需要几秒钟(最多 10 秒钟)。
希望对您有所帮助。查看这些链接了解更多信息:
http://docs.aws.amazon.com/streams/latest/dev/kinesis-record-processor-scaling.html
【讨论】: