【问题标题】:How to have fault tolerance on producer end with Kafka如何在生产者端使用 Kafka 进行容错
【发布时间】:2020-09-08 00:01:29
【问题描述】:

我是 Kafa 和数据摄取的新手。我知道 Kafka 是容错的,因为它将数据冗余地保存在多个节点上。但是,我不明白的是我们如何在源/生产者端实现容错。例如,如果我有 netcat 作为源,如下例所示。

nc -l [some_port] | ./bin/kafka-console-producer --broker-list [kafka_server]:9092 --topic [my_topic]

如果执行 netcat 的节点出现故障,生产者将无法推送消息。我在想是否有一种机制可以让 Kafka 自己提取输入,例如,如果一个节点上的 netcat 失败,另一个节点可以接管并开始使用 netcat 推送消息。

我的第二个问题是如何在 Flume 中实现这一点,因为它是基于拉取的架构。在这种情况下,即如果一个节点做 netcat 失败,Flume 会工作吗?

【问题讨论】:

  • 很想知道你是否得到了这个问题的答案!

标签: apache-kafka kafka-producer-api flume data-ingestion


【解决方案1】:

每个主题,都是一个特定的数据流(类似于数据库中的表)。主题,被分成 partitions(任意数量),其中分区中的每条消息都有一个增量 id,称为偏移量,如下所示。

分区 0:

+---+---+---+-----+
| 0 | 1 | 2 | ... |
+---+---+---+-----+

分区 1:

+---+---+---+---+----+
| 0 | 1 | 2 | 3 | .. |
+---+---+---+---+----+

现在一个 Kafka 集群由多个 brokers 组成。每个代理都有一个 ID 标识,并且可以包含某些主题分区。

2 个主题的示例(每个主题分别有 3 个和 2 个分区):

经纪人 1:

+-------------------+
|      Topic 1      |
|    Partition 0    |
|                   |
|                   |
|     Topic 2       |
|   Partition 1     |
+-------------------+

经纪人 2:

+-------------------+
|      Topic 1      |
|    Partition 2    |
|                   |
|                   |
|     Topic 2       |
|   Partition 0     |
+-------------------+

经纪人 3:

+-------------------+
|      Topic 1      |
|    Partition 1    |
|                   |
|                   |
|                   |
|                   |
+-------------------+

请注意,数据是分布式的(Broker 3 不保存 topic 2 的任何数据)。

主题,应该有一个replication-factor > 1(通常是 2 或 3),这样当一个代理关闭时,另一个可以提供主题的数据。例如,假设我们有一个主题有 2 个分区,replication-factor 设置为 2,如下所示:

经纪人 1:

+-------------------+
|      Topic 1      |
|    Partition 0    |
|                   |
|                   |
|                   |
|                   |
+-------------------+

经纪人 2:

+-------------------+
|      Topic 1      |
|    Partition 0    |
|                   |
|                   |
|     Topic 1       |
|   Partition 0     |
+-------------------+

经纪人 3:

+-------------------+
|      Topic 1      |
|    Partition 1    |
|                   |
|                   |
|                   |
|                   |
+-------------------+

现在假设 Broker 2 失败了。 代理 1 和 3 仍然可以为主题 1 提供数据。因此,replication-factor 为 3 始终是一个好主意,因为它允许一个代理被删除以进行维护,也允许另一个代理被删除出乎意料地被取下来。 因此,Apache-Kafka 提供了强大的持久性和容错保证。

关于领导者的注意事项: 在任何时候,只有一个代理可以成为分区的领导者,并且只有该领导者可以接收和提供该分区的数据。其余的代理只会同步数据(同步副本)。另请注意,当 replication-factor 设置为 1 时,leader 在代理失败时无法移动到其他位置。一般来说,当一个分区的所有副本都失败或下线时,leader 会自动设置为-1


话虽如此,只要您的生产者列出了集群中所有 Kafka 代理的地址 (bootstrap_servers),您应该没问题。即使一个代理关闭,您的生产者也会尝试将记录写入另一个代理。

最后,确保设置acks=all(虽然可能会影响吞吐量),以便所有同步副本确认他们收到了消息。

【讨论】:

  • 嗯,我理解你在这里描述的内容,但我的问题更多的是关于来源。这里我们在一个节点上运行 netcat。如果该节点发生故障会发生什么,另一个节点是否可以开始执行 netcat 向 Kafka 发送消息。 Kafka 是否为这种容错提供了任何东西?或者我们是否需要自己为源代码实现这一点。如果是这样,怎么做? Flume 在那里有用吗,因为它有一个基于拉的架构。
  • @MetallicPriest 这正是我的观点。只要您至少有 3 个代理,并且复制因子 = 3,您不必担心失去代理一段时间。
  • 嗯,你的意思是将 netcat 的输出推送给 Kakfa 生产者。这很清楚。但我的问题更多是关于源本身,即 netcat。如果执行 netcat 的节点失败,我们将不会向 Kafka 生产者提供任何消息。那么,基本上,我们如何才能在源头上实现容错呢?例如,我们能否让 Kafka 生产者自己运行 netcat 命令并将该输出推送到他们的主题/消息队列中?
  • 首先,你能澄清一下你所说的netcat是什么意思吗?您多次提到它,我不确定我是否理解上下文。
  • 好吧,netcat 是一个工具,在这个例子中,它用于从 Web 服务器读取流数据。因此,基本上 Kafka 生产者正在使用 netcat 进行馈送。但是您可以用任何来源替换 netcat,问题仍然存在。例如,对于 netcat 示例,Kafka 是否可以自行运行该命令。
猜你喜欢
  • 2016-01-15
  • 2020-04-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-05-05
  • 2019-04-27
  • 2023-03-22
相关资源
最近更新 更多