【问题标题】:Receiver less approach for spark-steaming with kinesis使用 kinesis 进行火花流的无接收器方法
【发布时间】:2017-04-25 14:50:28
【问题描述】:

对于使用 kafka 的 Spark 流,我们有 Directstream,它是无接收器的方法,并将 kafka 分区映射到 spark RDD 分区。目前我们有一个应用程序,我们在其中使用 Kafka Direct 方法并在 RDBMS 中维护我们的偏移量,

我们有类似的 Kinesis 吗?当我阅读 spark-Kinesis 集成的文档时,感觉检查点存在差异。以下是我的一些问题

  1. 使用 kinesis 流式处理是否将 kinesis 分片映射到 RDD 分区?如果我在传入的 RDD 上使用 forEachPartition,我可以在分片级别保持有序处理吗?
  2. 从文档中解释了 kinesis 在 dynamoDB 中维护单独的检查点?我们不能忽略它并使用我们自己的偏移管理吗?
  3. 在 KinesisUtils.createStream api 中,我看到对于 [初始位置] 变量,它只需要 LATEST 或 TRIM_HORIZON。在那种情况下,我怎么能不能像我在 kafka 案例中提供的那样提供要偏移的分片映射?

如果我们的应用程序是幂等的,我们如何才能得到恰好一次处理?

【问题讨论】:

    标签: apache-spark spark-streaming amazon-kinesis


    【解决方案1】:

    使用 kinesis 流式处理是否将 kinesis 分片映射到 RDD 分区?

    不,Kinesis 分片和 RDD 分区之间不存在 1:1 映射,如 documentation 中所述:

    在输入 DStream 处理期间,Kinesis 流分片的数量与跨 Spark 集群创建的 RDD 分区/分片的数量之间没有关联。这是 2 个独立的分区方案


    如果我在传入的 RDD 上使用 forEachPartition,我能否在分片级别保持有序处理?

    每个创建的分区,内部都保持顺序(不确定是否有帮助):

    Kinesis 数据处理按分区排序,并且每条消息至少发生一次。


    从文档中解释了 kinesis 在 dynamoDB 中维护单独的检查点?我们不能忽略它并使用我们自己的偏移管理吗?

    不,您受到使用 DyanmoDB 作为后备存储的 Kinesis 客户端实现的约束。

    在 KinesisUtils.createStream api 中,我看到对于 [初始位置] 变量,它只需要 LATEST 或 TRIM_HORIZON。在那种情况下,我怎么不能像我在 kafka 案例中提供的那样提供要偏移的分片映射?

    没有。没有提供等效的 Kafka 偏移量。

    如您所见,Kinesis API 的当前实现限制了您。如果您需要偏移存储和恢复的灵活性,并希望实现恰好一次的语义,那么也可以考虑使用 Kafka 来实现此解决方案。

    【讨论】:

    • 感谢您的回复。在那种情况下,kinesis 远远超出了 kafka。这对我的云迁移来说是个大问题。我有两个关键应用程序只使用一次火花流。
    猜你喜欢
    • 1970-01-01
    • 2017-04-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-11
    相关资源
    最近更新 更多