【问题标题】:consumer reading __consumer_offsets delivers unreadable message消费者阅读 __consumer_offsets 传递不可读的消息
【发布时间】:2019-02-16 02:10:49
【问题描述】:

我正在尝试从 __consumer_offsets 主题中消费,因为这似乎是检索有关消费者的 kafka 指标(如消息滞后等)的最简单方法。理想的方法是从 jmx 访问它,但想先尝试一下回来似乎是加密的或不可读的形式。也尝试添加 stringDeserializer 属性。有没有人对如何纠正这个有任何建议?再次引用 this 是

的副本

duplicate consumer_offset

没有帮助,因为它没有引用我的问题,即在 java 中将消息作为字符串读取。还更新了代码以尝试使用 kafka.client 消费者的 consumerRecord。

consumerProps.put("exclude.internal.topics",  false);
consumerProps.put("group.id" , groupId);
consumerProps.put("zookeeper.connect", zooKeeper);


consumerProps.put("key.deserializer",
  "org.apache.kafka.common.serialization.StringDeserializer");  
consumerProps.put("value.deserializer",
  "org.apache.kafka.common.serialization.StringDeserializer");

ConsumerConfig consumerConfig = new ConsumerConfig(consumerProps);
ConsumerConnector consumer = 
kafka.consumer.Consumer.createJavaConsumerConnector(
       consumerConfig);

Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
topicCountMap.put(topic, new Integer(1));
Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = 
   consumer.createMessageStreams(topicCountMap);
List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);

for (KafkaStream stream : streams) {

    ConsumerIterator<byte[], byte[]> it = stream.iterator();

    //errorReporting("...CONSUMER-KAFKA CONNECTION SUCCESSFUL!");

    while (it.hasNext()) {
         try {

             String mesg = new String(it.next().message());
             System.out.println( mesg);

代码更改:

try {       
    // errorReporting("CONSUMER-KAFKA CONNECTION INITIATING...");   
    Properties consumerProps = new Properties();
    consumerProps.put("exclude.internal.topics",  false);
    consumerProps.put("group.id" , "test");
    consumerProps.put("bootstrap.servers", servers);
    consumerProps.put("key.deserializer","org.apache.kafka.common.serialization.StringDeserializer");  
    consumerProps.put("value.deserializer","org.apache.kafka.common.serialization.StringDeserializer");

    //ConsumerConfig consumerConfig = new ConsumerConfig(consumerProps);
    //ConsumerConnector consumer = kafka.consumer.Consumer.createJavaConsumerConnector(
    //       consumerConfig);

    //Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
    //topicCountMap.put(topic, new Integer(1));
    //Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumer.createMessageStreams(topicCountMap);
    //List<KafkaStream<byte[], byte[]>> streams = consumerMap.get(topic);

    KafkaConsumer<String, String> kconsumer = new KafkaConsumer<>(consumerProps); 
    kconsumer.subscribe(Arrays.asList(topic)); 

    try {
        while (true) {
            ConsumerRecords<String, String> records = kconsumer.poll(10);

            for (ConsumerRecord<String, String> record : records)

                System.out.println(record.offset() + ": " + record.value());
        }
    } finally {
          kconsumer.close();
    }    

下面是消息的快照;在图片底部:

【问题讨论】:

    标签: java apache-kafka


    【解决方案1】:

    虽然可以直接从__consumer_offsets 主题读取,但这不是推荐或最简单的方法。

    如果可以使用 Kafka 2.0,最好使用 AdminClient API 来描述组:


    如果您绝对想直接从__consumer_offset 中读取,则需要对记录进行解码以使其可读。这可以使用GroupMetadataManager 类来完成:

    您链接的问题中的这个answer 包含执行所有这些的骨架代码。

    还请注意,您不应将记录反序列化为字符串,而是将它们保留为原始字节,以便这些方法能够正确解码它们。

    【讨论】:

    • 感谢您的回复。我还没有解决这个问题,但你的详细解释是我从任何地方得到的最好的回应。欣赏它。
    • 我不确定为什么要重新提出这个问题。在写完您链接到的答案后,我已将其作为副本关闭。我也可以在那里引用 AdminClient 方法
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-13
    • 2017-11-09
    • 1970-01-01
    相关资源
    最近更新 更多