【问题标题】:Can flume, spark, storm read from Kafka in round-robin way?可以循环方式从卡夫卡读取水槽、火花、风暴吗?
【发布时间】:2016-02-04 03:37:49
【问题描述】:

我有 3 个分区:0、1、2。所以消息可以分为 0、1、2。

例如:

分区 0 中的 1 条消息:0

分区 1 中的 3 条消息:111

分区 2 中有 2 条消息:22

如何让消费者按照012x12x1x的顺序消费消息(x表示当时没有消息)。消费消息的顺序如下:012121。 我想在 C++ 和 Python 中都这样做。查看现有客户端,消息可以以循环方式生成,但不能以循环方式消费。

有什么想法吗?

Kafka 消费者配置中有 partition.assignment.strategy (http://kafka.apache.org/documentation.html#consumerconfigs)。我正在寻找一些实现此配置的工具(例如flume、spark、storm)来从kafka读取数据,重新排序,然后再次写入kafka。继续上面的例子。重新排序的消息如下所示: 012121 (012x12x1x)

更新

现在,我可以在 C++ Kafka 客户端 (https://github.com/edenhill/librdkafka) 中执行此操作。

    for(int i = 0; i < 2; i++)
    {
        RdKafka::Message *msg = m_consumer->consume(m_topic, i, 1000);
        // Do something about msg here...
    }

输出:

Reading from 1=>4953---1---
Reading from 0=>46164---0---
Reading from 1=>4954---1---
Reading from 0=>46165---0---
Reading from 1=>4955---1---
Reading from 0=>46166---0---
Reading from 1=>4956---1---
Reading from 0=>46167---0---
Reading from 1=>4957---1---
Reading from 0=>46168---0---

【问题讨论】:

  • 我的问题已重新编辑。谢谢。

标签: apache-kafka spark-streaming flume apache-storm round-robin


【解决方案1】:

在 Java 中,您可以使用 high level consumer

如果您从同一进程中的所有 3 个分区消费,使用相同的 group.id,您可以创建 3 个充当消费者的工作线程,并以循环方式在它们之间循环。

我了解到您明确提到 Python 和 C++ 作为您的目标语言。
根据我自己在 Python 中的经验,没有与 Java 中提供的高级消费者等效的东西。

因此,您可以以某种方式包装它,例如创建一个充当服务器的新 Java 进程,并通过该进程使用消息,或者您可以尝试将高级消费者移植到 Python 和/或 C++。

另一种选择是使用每个分区,将数据写入其他介质,如 MySQL(或任何其他 RDBMS),并使用 SQL 使用消息来执行循环,最终将它们从相关表中删除。

无论如何,我建议您重新考虑将 Kafka 作为您的传输层。您的需求(以循环方式消费)与 Kafka 的核心设计/架构不一致。原因如下:
只要分区的数量小于消费者机器中的可用内核 - 就可以像我上面描述的那样以循环方式使用那台(!)机器上的消息。

不过,Kafka 也是为(消费者的)分布式负载而设计的。这就是为单个主题创建分区的动机。可扩展性是 Kafka 中的一个关键概念。

因此,如果您的分区比消费者机器中的可用内核数多,您可能无法正确迭代分区并以循环方式使用消息。

示例:假设您在特定主题中有 20 个分区。这个话题总是会产生很多消息。现在假设您有一台具有 4 个 CPU 内核的机器可以从该主题中使用。按照设计,您一次最多可以使用 4 个分区。可以通过在内存中缓冲大量消息或其他一些机制(这可能导致其他问题,如延迟)来实现将其转换为循环消耗。那时我们假设所有分区始终可用,没有影响 Kafka 的网络问题或磁盘问题。这就是为什么我认为将分区“重新加入”到其他一些消息流中并不是一件小事。

【讨论】:

  • Kafka 消费者配置中有 partition.assignment.strategy (kafka.apache.org/documentation.html#consumerconfigs)。我正在寻找一些实现此配置的工具(例如 kafka logstash 插件)以从 kafka 读取数据,重新排序它们,然后再次写入 kafka。你有什么建议吗?
  • [a] 我不知道这些工具; [b] 我知道 Logstash Kafka 插件,但它的消费策略非常简单; [c] 我真的认为您试图实现的目标与 Kafka 的核心设计决策相冲突。 [d] 因此,我谦虚地建议您使用来自 Kafka 的数据,使用队列或 RDBMS 等其他介质,根据需要重新排序,最后在不同的主题中推回 Kafka。
  • 可以自己写kafka输入插件吗?在 Java Kafka 客户端中,支持循环消费。所以,这个客户端可以用于我自己的插件。
  • 尝试在水槽、火花或风暴中使用。现在,我正在谷歌上搜索。
  • Kafka 的核心设计/架构中没有禁止轮询消费的内容。消费者应用程序可以根据需要在每次迭代中简单地使用来自 N 个分区的 M 条消息。这方面的一个例子是新的 Kafka 流 API,它支持诸如 join(partitions..) 之类的东西。
猜你喜欢
  • 2016-08-03
  • 2018-09-15
  • 2019-02-18
  • 1970-01-01
  • 2018-04-25
  • 2018-02-24
  • 2023-03-19
  • 1970-01-01
  • 2018-08-13
相关资源
最近更新 更多