【发布时间】: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