【问题标题】:How to distribute kafka load into openshift pods using kafka partitions如何使用 kafka 分区将 kafka 负载分配到 openshift pod
【发布时间】:2021-03-10 20:38:36
【问题描述】:

我有一个 Spring Boot 应用程序,我想将 Kafka 主题的负载分配到 3 个 open-shift pod。我有以下示例,我可以在三个不同线程上监听 3 个 Kafka 分区,这个 Spring Boot 应用程序将加载到一个 openshift pod 中。但是我希望能够从一个 pod 上的一个 Kafka 分区监听,所以当我在 open-shift 上加载 3 个 pod 时,每个 pod 都会从一个 Kafka 分区监听。这将允许我将应用程序扩展到 N 个 pod 上的 N 个分区。我不确定这是否可行,或者是否需要使用不同的方法。谢谢

public class DepAcctInqConsumerController {
   
   private static final Logger LOGGER = LoggerFactory.getLogger(DepAcctInqConsumerController.class);

   @Value("${kafka.topic.acct-info.request}")
   private String requestTopic;

  @KafkaListener(id = "id-0",containerFactory = "requestReplyListenerContainerFactory",
          topicPartitions = { @TopicPartition(topic = "${kafka.topic.acct-info.request}", partitions = "0" )})
  public Message<?> listenPartition0(InGetAccountInfo accountInfo, @Header(KafkaHeaders.REPLY_TOPIC) byte[] replyTo,
                           @Header(KafkaHeaders.CORRELATION_ID) byte[] correlation,@Header(KafkaHeaders.RECEIVED_PARTITION_ID) int id) {

    try {

      LOGGER.info("Received request for partition id = " +  id);

      AccountInquiryDto accountInfoDto  = getAccountInquiryDto(accountInfo);

      return MessageBuilder.withPayload(accountInfoDto)
              .setHeader(KafkaHeaders.TOPIC, replyTo)
              .setHeader(KafkaHeaders.RECEIVED_PARTITION_ID, id)
              .setHeader(KafkaHeaders.CORRELATION_ID, correlation)
              .build();


    } catch (Exception e) {
      LOGGER.error(e.toString(),e);
    }

    return null;
  }


  @KafkaListener(id = "id-1",containerFactory = "requestReplyListenerContainerFactory",
          topicPartitions = { @TopicPartition(topic = "${kafka.topic.acct-info.request}", partitions = "#{@finder.partitions(${kafka.topic.acct-info.request)}" )})
  public Message<?> listenPartition1(InGetAccountInfo accountInfo, @Header(KafkaHeaders.REPLY_TOPIC) byte[] replyTo,
                           @Header(KafkaHeaders.CORRELATION_ID) byte[] correlation,@Header(KafkaHeaders.RECEIVED_PARTITION_ID) int id) {

    try {

      LOGGER.info("Received request for partition id = " +  id);

      AccountInquiryDto accountInfoDto  = getAccountInquiryDto(accountInfo);

      return MessageBuilder.withPayload(accountInfoDto)
              .setHeader(KafkaHeaders.TOPIC, replyTo)
              .setHeader(KafkaHeaders.RECEIVED_PARTITION_ID, id)
              .setHeader(KafkaHeaders.CORRELATION_ID, correlation)
              .build();


    } catch (Exception e) {
      LOGGER.error(e.toString(),e);
    }

    return null;
  }

  @KafkaListener(id = "id-2",containerFactory = "requestReplyListenerContainerFactory",
          topicPartitions = { @TopicPartition(topic = "${kafka.topic.acct-info.request}", partitions =  "2" )})
  public Message<?> listenPartition2(InGetAccountInfo accountInfo, @Header(KafkaHeaders.REPLY_TOPIC) byte[] replyTo,
                           @Header(KafkaHeaders.CORRELATION_ID) byte[] correlation, @Header(KafkaHeaders.RECEIVED_PARTITION_ID) int id) {

    try {

      LOGGER.info("Received request for partition id = " +  id);

      AccountInquiryDto accountInfoDto  = getAccountInquiryDto(accountInfo);

      return MessageBuilder.withPayload(accountInfoDto)
              .setHeader(KafkaHeaders.TOPIC, replyTo)
              .setHeader(KafkaHeaders.RECEIVED_PARTITION_ID, id)
              .setHeader(KafkaHeaders.CORRELATION_ID, correlation)
              .build();


    } catch (Exception e) {
      LOGGER.error(e.toString(),e);
    }

    return null;
  }

【问题讨论】:

    标签: java spring-boot apache-kafka openshift


    【解决方案1】:

    我们不需要为每个分区设置多个 kafka 监听器。我们只需要一位听众。

    • 如果您正在运行单个 pod,则来自所有三个分区的消息都将由该单个 pod 使用,
    • 如果您运行超过 1 个 pod,分区将分布在 pod 之间。
    • 我们可以运行与没有分区一样多的 pod。

    所有 Pod 必须使用相同的消费者组名称。

    这就是我们所需要的。

    @KafkaListener(topics = "${kafka.topic.acct-info.request}")
    public void receive(ConsumerRecord<String, String> record)
    

    【讨论】:

      猜你喜欢
      • 2017-03-25
      • 2019-10-08
      • 2018-10-08
      • 2021-01-26
      • 1970-01-01
      • 1970-01-01
      • 2018-06-21
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多