【问题标题】:Run process on any one Pod in the Kubernetes deployment在 Kubernetes 部署中的任意一个 Pod 上运行进程
【发布时间】:2021-03-25 02:35:46
【问题描述】:

我有一个在多个 pod 上运行并在流量增加时横向扩展的应用程序。该应用程序的功能之一是“它从 Kafka 中挑选消息并触发电子邮件”。

当多个 pod 运行时,它们都会触发电子邮件,因为所有 pod 都应该按照设计接收消息。

如何限制电子邮件功能一次在任何一个 pod 上工作?

集群 - EKS ,编程语言 - Scala AKKA

【问题讨论】:

  • 当消费者处理消息时,消息不会从其主题中删除。相反,消费者可以从多种方式中选择让 Kafka 知道哪些消息已被处理。这个过程称为提交偏移量。因此,一旦您从 list 获取消息,您需要立即提交 offset 以便其他消费者将收到以下消息
  • “限制电子邮件功能一次只能在任何一个 pod 上工作”到底是什么意思?
  • @LeviRamsey 表示所有 pod 上都有电子邮件触发功能,但只有一个 pod 应该执行它。
  • @Pranav 如果只有一个 pod 触发了电子邮件,其他 pod 会做什么?使用消费者组可以保证一些邮件被 pod1、podx 等触发,从而分担工作量。这样,同一封邮件就不会被多个 Pod 触发。
  • @JavaTechnical 问题在于 Pod 还提供有关与前端应用程序的 Web 套接字连接的实时数据,并且前端应用程序连接到任何一个 Pod。所以来自 kafka 的数据应该被所有的 Pod 获取。实际上,相同的数据也是通过电子邮件以合并的形式推送的,

标签: kubernetes apache-kafka akka


【解决方案1】:

如何限制电子邮件功能在任何一个 一个豆荚?

简而言之:对所有触发电子邮件的 pod 使用相同的消费者组。 通常,工作负载会根据它们所做的工作分为几组。同一组的成员彼此分担工作量。

您当然可以将 bootstrap.servers 等 Kafka 消费者配置提供给您的 pod。在该配置中,将名称为 group.id 的属性赋予某个值,例如 email-trigger-group,然后将按照您的预期共享工作负载。

您可以为触发电子邮件的 pod 使用 标签。您可以为您的消费者 group.id 为您的所有 pod 使用相同的标签值。


我们可以把问题分成两个子问题:

1.触发电子邮件

此工作负载可以由组中的多个使用者共享。

2。向前端应答请求

对整个主题(所有分区)使用手动consumer.assign()。

前端将指定 时间戳,它希望从中获取新消息,即时间戳为 > 的消息将从主题的所有分区中检索此时间戳。为此,请使用consumer.offsetsForTimes() 获取时间戳、轮询并将消息作为响应发送。

List<TopicPartition> topicPartitions = consumer.partitionsFor("your_topic").stream().map(partitionInfo -> new TopicPartition(partitionInfo.topic(), partitionInfo.partition()).toList();
consumer.assign(topicPartitions);

// Populate the map
Map<TopicPartition, Long> partitionTimestamp = new LinkedHashMap<>();

// Add the same timestamp received from frontend for all partitions
topicPartitions.forEach(topicPartition -> partitionTimestamp.put(topicPartition, timestampFromFrontend));

// Get the offsets and seek
Map<TopicPartition,OffsetAndTimestamp> offsetsForTimes = consumer.offsetsForTimes(offsetsForTimestamp);

// Seek to the offsets
offsetsForTimes.forEach( (tp, oft) -> consumer.seek(tp, oft.offset()) );

// Poll and return
consumer.poll();

【讨论】:

  • 问题在于 Pod 还提供了与前端应用程序的 Web 套接字连接的实时数据,并且前端应用程序连接到任何一个 Pod。所以来自 kafka 的数据应该被所有 pod 获取。
  • @Pranav 因此,前端连接到的任何 pod 都必须提供来自所有分区的消息,直到那时(无论它被分配到哪个分区)。不是吗?
  • 使用consumer.assign() 而不是subscribe() 来回答前端请求,并要求前端给出它最后一条消费消息的时间戳,它需要新消息。只需使用offsetForTimes() 查找由该时间戳标识的所有分区的偏移量,轮询并返回响应。
【解决方案2】:

如果您使用的是 Kafka,则可以使用分区。

每个主题可以有多个分区。这些分区在消费者之间共享。

例如:

Email Topic: Partitions[0,1,2,3,4,5]

Email Consumer Group:
   Consumer 1: Partitions[0,3]
   Consumer 2: Partitions[1,4]
   Consumer 3: Partitions[2,5]

On Scale Up Event:
   Consumer 1: Partitions[0]
   Consumer 2: Partitions[1]
   Consumer 3: Partitions[2,5]
   Consumer 4: Partitions[3]
   Consumer 5: Partitions[4]

On Scale Down Event:
   Consumer 1: Partitions[0,2,4]
   Consumer 2: Partitions[1,3,5]

这样只有特定组中的一个消费者会消费消息。

【讨论】:

  • 我需要所有 pod 来接收所有数据,因为除了电子邮件部分之外,它们还有一些其他功能
  • 电子邮件组已分区,因此它们共享负载。您可以将其他组添加到同一主题。每个组将收到相同的主题数据。休息取决于你想如何分区。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2019-04-23
  • 1970-01-01
  • 2019-02-03
  • 2020-08-05
  • 2020-07-20
  • 2023-04-03
  • 2023-01-10
相关资源
最近更新 更多