【问题标题】:how to do Kafka syncing between two consumer groups如何在两个消费者组之间进行 Kafka 同步
【发布时间】:2019-12-21 17:57:58
【问题描述】:

在我的要求中,我有两个消费者组,一个组(主)只是获取数据并发送到其他服务器,如果发送到其他服务器失败,那么我需要重新加入(启动)其他消费者组(处理失败) .

在这种情况下,主组继续读取和重试,它会继续相同,当消息发送成功时需要通知其他消费者组(处理失败)。现在失败的处理应该从第一次失败到最后一次失败的地方开始发送。

public void StartMainstreamHandler() {
       StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> userStream = builder.stream("usertopic",Consumed.with(Serdes.String(), Serdes.String()));
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "main-streams-userstream");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "ALL my bootstrap servers);
        props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4);
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "500");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        //consumer_timeout_ms
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 2000);

        props.put("state.dir","/tmp/kafka/stat));


   userStream.peek((key,value)->System.out.println("key :"+key+" value :"+value));
   /* Send Data to other Server if Failed call other consumer */

KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), props);


kafkaStreams.setUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() {
            @Override
            public void uncaughtException(Thread t, Throwable e) {

                logger.error("Thread Name :" + t.getName() + " Error while processing:", e);
            }
        });


        kafkaStreams.cleanUp();
        kafkaStreams.start();
    }

其他消费者

public void StartFailstreamHandler() {
       StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> userStream = builder.stream("usertopic",Consumed.with(Serdes.String(), Serdes.String()));
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "failed-streams-userstream");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "ALL my bootstrap servers);
        props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 4);
        props.put("enable.auto.commit", "true");
        props.put("auto.commit.interval.ms", "500");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        //consumer_timeout_ms
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 2000);

        props.put("state.dir","/tmp/kafka/stat));

     Wait('till get notfication from other consumer" ){
     userStream.peek((key,value)->System.out.println("key :"+key+" value :"+value));
     /* start sending */
     /* how to break when it is reached last offest */ 
   }

KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), props);


kafkaStreams.setUncaughtExceptionHandler(new Thread.UncaughtExceptionHandler() {
            @Override
            public void uncaughtException(Thread t, Throwable e) {

                logger.error("Thread Name :" + t.getName() + " Error while processing:", e);
            }
        });


        kafkaStreams.cleanUp();
        kafkaStreams.start();
    }

现在如何知道和同步第二个消费者的偏移细节,以在最后失败时完全停止(最后失败发生在主要消费者)

【问题讨论】:

  • 所以你基本上用 Kafka Streams 重写了 MirrorMaker?
  • 实际目标是保留一个消费者组未成功处理的消息,而保留的消息应该由另一个消费者组处理。
  • 所以你正在制作一个死信主题?那你为什么需要同步组呢?偏移量完全不同
  • 我在想,一旦消费者被读取,我们需要两个消费者组来将消息保留在主题中,后来知道我可以使用保留期medium.com/@werneckpaiva/… .. 现在的问题是如何知道关于两个正在读取最新数据和最早读取其他数据的消费者(失败的时间戳),我想在另一个消费者到达成功处理数据的主要消费者时停止它
  • 您需要在外部存储该信息。消费群体无法开箱即用地相互协调

标签: java apache-kafka apache-kafka-streams


【解决方案1】:

仅使用消息在集群之间同步消费者组是不可能的,因为消费者无法找到特定的“复制起点”。

您必须将额外的元数据存储到一边,例如复制开始的时间戳,可能将该信息嵌入到记录标头中(假设您的 Kafka 版本支持它}。否则,您或多或少是在盲目地复制数据以最低限度的交付保证(因此无论如何使用 MirrorMaker 可能会更好)。

MirrorMaker 2 或 Confluent Replicator 是目前仅有的两个可用于“同步”消费者组的选项,以进行主动-主动双集群设置

【讨论】:

  • MM2 使两个 DC 电平同步。即使对于同一 DC 中的两个消费者组,这是唯一的方法吗?
  • 有什么建议吗?
  • 我不明白关于“跨同一个 DC 中的两个消费者组”的问题。为什么要在同一个 DC 中同步,而不是让所有消费者都属于一个组?
  • 正如我所提到的,我的服务需要读取数据并发送到一些第三方服务器,所以有时第三方服务在他们的服务启动后出现故障,我需要发送最新的记录和同时失败的记录。所以我在消费者或群体之间进行了这种节日同步。
  • 您不能在数据库中存储偏移量吗?如果没有被另一个 kafka 消费者实际使用,为什么他们只需要在一个消费者组中?另外,我认为 kafka 流无论如何都不会公开偏移信息
猜你喜欢
  • 1970-01-01
  • 2020-07-08
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-01-07
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多