【问题标题】:Kafka - Dynamic / Arbitrary PartitioningKafka - 动态/任意分区
【发布时间】:2015-05-08 23:50:51
【问题描述】:

我正在为 Kafka 主题构建消费者服务。每条消息都包含一个 url,我的服务将向其发出 http 请求。每个消息/url 完全独立于其他消息/url。

我担心的问题是如何处理长时间运行的请求。某些 http 请求可能需要 50 多分钟才能返回响应。在那段时间里,我不想保留任何其他消息。

并行化此操作的最佳方法是什么?

我知道 Kafka 的并行方法是创建分区。但是,根据我的阅读,当我真的想要无限或动态的分区数时,您似乎需要预先定义分区数(理想情况下,每条消息都会即时创建自己的分区)

例如,假设我创建了 1,000 个分区。如果为我的主题生成了 1,001 多条消息,则会发出前 1,000 个请求,但之后的每条消息都将排队,直到该分区中的前一个请求完成。

我曾考虑过使 http 请求异步,但在确定要提交的偏移量时似乎遇到了问题。

例如,在单个分区上,我可以让消费者读取第一条消息并发出异步请求。它提供了一个回调函数,该函数将该偏移量提交给 Kafka。在该请求等待时,我的消费者会读取下一条消息并发出另一个异步请求。如果该请求在第一个请求之前完成,它将提交该偏移量。现在,如果第一个请求由于某种原因失败或我的消费者进程死亡会发生什么?如果我已经提交了更高的偏移量,听起来这意味着我的第一条消息将永远不会被重新处理,这不是我想要的。

在使用 Kafka 进行长时间运行的异步消息处理时,我显然遗漏了一些东西。有没有人遇到过类似的问题或对如何最好地解决这个问题有想法?提前感谢您抽出宝贵时间阅读本文。

【问题讨论】:

    标签: asynchronous apache-kafka job-scheduling


    【解决方案1】:

    您应该查看 Apache Storm 的消费者处理部分,并将消息存储和检索留给 Kafka。您所描述的是大数据中非常常见的用例(尽管 50 多分钟的事情有点极端)。简而言之,您将为您的主题设置少量分区,并让 Storm 流处理扩展实际发出 http 请求的组件(Storm 中的“螺栓”)的数量。单个 spout(一种从外部源读取数据的风暴组件)可以从 Kafka 主题读取消息并将它们流式传输到处理螺栓。

    我已经在 github 上发布了 an open source example 如何编写 Storm/Kafka 应用程序。

    对此答案的一些后续想法:

    1) 虽然我认为 Storm 是正确的平台方法,但您没有理由不通过编写一个执行 http 调用的 Runnable 然后编写更多代码来让单个 Kafka 消费者读取消息来自行开发并使用可运行的多线程实例处理它们。所需的管理代码有点有趣,但可能比从头开始学习 Storm 更容易编写。因此,您可以通过在更多线程上添加更多 Runnable 实例来进行扩展。

    2)无论你使用Storm还是你自己的多线程解决方案,你仍然会遇到如何管理Kafka中的偏移量的问题。简短的回答是您必须自己进行复杂的偏移管理。您不仅需要保留从 Kafka 读取的最后一条消息的偏移量,而且还必须保留和管理当前正在处理的正在处理的消息列表。这样,如果您的应用程序出现故障,您就可以知道正在处理哪些消息,并且可以在启动备份时检索和重新处理它们。基本的 Kafka 偏移持久性不支持这种更复杂的需求,但无论如何它只是为了方便更简单的用例。您可以在任何您喜欢的地方(动物园管理员、文件系统或任何数据库)保存您的偏移信息。

    【讨论】:

    • 谢谢,克里斯!我有研究流处理(Storm、Samza 等)的长期计划,但鉴于当前的截止日期,这个当前的问题需要一个短期的解决方案。对保留解决方案有任何想法吗?
    • 刚刚添加了那些 cmets
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-02-04
    • 2018-10-24
    • 1970-01-01
    • 1970-01-01
    • 2017-09-13
    • 2022-01-01
    • 1970-01-01
    相关资源
    最近更新 更多