【发布时间】: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