【问题标题】:IRecordsProcessor processRecords not invoked by KCL libraryKCL 库未调用 IRecordsProcessor processRecords
【发布时间】:2016-10-12 15:14:31
【问题描述】:

我们多次进行了以下测试,但我们很难找到一个合理的解释来解释为什么会发生这种情况:

  • 我们创建一个消费者,等待它准备好
  • 我们在消费者正在收听的流上发布两条记录

我们有一个分片,有时消费者不会收到有关记录的通知。我们使用不同的 workerId,但具有相同 ApplicationName 的应用程序可能会窃取记录。

KCL 消费者从未获得刚刚发布的记录的原因是什么?

【问题讨论】:

    标签: amazon-kinesis


    【解决方案1】:

    我的代码的问题是,一旦您第一次在 KCL 中创建应用程序,Dynamo 中用于检查点的匹配条目会存储分片迭代器的最后位置。

    当您重新启动应用程序时,分片迭代器可能会比当前时间滞后很多,并且在您到达发布记录的当前位置之前可能需要进行大量迭代。

    因此,出于测试目的,您可能需要每次都创建一个新应用并每晚清理您的发电机表

    【讨论】:

    • 您正在使用 KCL,它使用 DynamoDB 表处理检查点。当工作人员重新启动时,它开始从最后一个检查点读取事件。在 KCL 中无法更改此行为。
    • 您可以直接从分片读取流中发布的新事件,但这意味着放弃使用 KCL 并使用分片。当您从分片读取时,可以指定迭代器类型,例如最新的。请在此处查看更多详细信息:docs.aws.amazon.com/kinesis/latest/APIReference/…
    • @prisco.napoli 的重点是它很长时间没有读取任何内容(它说没有可用的记录)
    • 似乎更像是您的应用程序逻辑中的一个错误,但不看代码很难说。其他工作人员是否正确接收消息?
    • Kinesis 客户端(而不是我的应用程序)记录没有可用的记录。我刚刚降低了记录器级别:14:18:04.978 [ForkJoinPool-1-worker-3] DEBUG c.a.s.k.c.lib.worker.ProcessTask - Kinesis 没有返回任何分片 shardId-000000000000 的记录
    【解决方案2】:

    您似乎在每个应用程序名称中使用了许多工作人员,但您的流只有一个分片。每个分片只能有一个工作人员处于活动状态并检索数据,其他所有工作人员都将处于空闲状态。

    您应确保工作人员的数量不超过分片的数量(或仅增加几个以提高弹性)。

    使用 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-implementation-app-java.html#kcl-java-worker

    http://docs.aws.amazon.com/streams/latest/dev/kinesis-record-processor-scaling.html

    【讨论】:

    • 所以在测试时我应该总是使用不同的应用程序名称?我试过了,但没有用。当我继续使用发电机时,我找不到我的表与流名称的关系
    • 如果您想测试多个使用同一分区中数据的工作人员,您可以为每个工作人员使用不同的应用程序名称。否则,将始终有一名工作人员处于活动状态。
    • 谢谢,如果一个应用程序从两个流中消费,我需要两个单独的应用程序名称吗?
    • 是的,如果您想要一个使用来自不同 kinesis 流的数据的应用程序,您需要为每个流定义至少一个工作器。每个工作人员都有不同的应用程序名称。“应用程序名称”只是一个标签,用于将在同一流上工作的许多工作人员组合在一起。
    • 我在 API 中检查得更好,我怀疑对于使用来自不同流的数据的工作人员也可能使用相同的“应用程序名称”。我从未测试过这种情况,但似乎 API 允许这样做。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-04-15
    • 1970-01-01
    • 1970-01-01
    • 2021-01-19
    • 2020-03-08
    • 1970-01-01
    相关资源
    最近更新 更多