【问题标题】:Kafka consumer in java not consuming messagesjava中的Kafka消费者不消费消息
【发布时间】:2015-04-30 10:46:33
【问题描述】:

我正在尝试让 kafka 消费者获取生成并发布到 Java 主题的消息。我的消费者如下。

consumer.java

import java.io.UnsupportedEncodingException;
import java.nio.ByteBuffer;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;

import kafka.consumer.Consumer;
import kafka.consumer.ConsumerConfig;
import kafka.consumer.ConsumerIterator;
import kafka.consumer.KafkaStream;
import kafka.javaapi.consumer.ConsumerConnector;
import kafka.javaapi.message.ByteBufferMessageSet;
import kafka.message.MessageAndOffset;



public class KafkaConsumer extends  Thread {
    final static String clientId = "SimpleConsumerDemoClient";
    final static String TOPIC = " AATest";
    ConsumerConnector consumerConnector;


    public static void main(String[] argv) throws UnsupportedEncodingException {
        KafkaConsumer KafkaConsumer = new KafkaConsumer();
        KafkaConsumer.start();
    }

    public KafkaConsumer(){
        Properties properties = new Properties();
        properties.put("zookeeper.connect","10.200.208.59:2181");
        properties.put("group.id","test-group");      
        ConsumerConfig consumerConfig = new ConsumerConfig(properties);
        consumerConnector = Consumer.createJavaConsumerConnector(consumerConfig);
    }

    @Override
    public void run() {
        Map<String, Integer> topicCountMap = new HashMap<String, Integer>();
        topicCountMap.put(TOPIC, new Integer(1));
        Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = consumerConnector.createMessageStreams(topicCountMap);
        KafkaStream<byte[], byte[]> stream =  consumerMap.get(TOPIC).get(0);
        System.out.println(stream);
        ConsumerIterator<byte[], byte[]> it = stream.iterator();
        while(it.hasNext())
            System.out.println("from it");
            System.out.println(new String(it.next().message()));

    }

    private static void printMessages(ByteBufferMessageSet messageSet) throws UnsupportedEncodingException {
        for(MessageAndOffset messageAndOffset: messageSet) {
            ByteBuffer payload = messageAndOffset.message().payload();
            byte[] bytes = new byte[payload.limit()];
            payload.get(bytes);
            System.out.println(new String(bytes, "UTF-8"));
        }
    }
}

当我运行上面的代码时,我在控制台中什么也得不到,屏幕后面的 java 生产者程序在“AATest”主题下连续发布数据。此外,在 Zookeeper 控制台中,当我尝试运行上述 consumer.java 时,我得到以下几行

[2015-04-30 15:57:31,284] INFO Accepted socket connection from /10.200.208.59:51780 (org.apache.zookeeper.
server.NIOServerCnxnFactory)
[2015-04-30 15:57:31,284] INFO Client attempting to establish new session at /10.200.208.59:51780 (org.apa
che.zookeeper.server.ZooKeeperServer)
[2015-04-30 15:57:31,315] INFO Established session 0x14d09cebce30007 with negotiated timeout 6000 for clie
nt /10.200.208.59:51780 (org.apache.zookeeper.server.ZooKeeperServer)

此外,当我运行一个指向 AATest 主题的单独控制台消费者时,我会获取生产者生成的所有数据到该主题。

消费者和代理都在同一台机器上,而生产者在不同的机器上。这实际上类似于this question。但是经历它对我有帮助。请帮帮我。

【问题讨论】:

  • 您尝试添加props.put("auto.offset.reset", "smallest"); 吗?
  • 是的,我试过了,但得到了相同的结果,消费者端没有数据。 (抱歉延迟回复)
  • 非常抱歉..问题是因为主题名称前的拼写错误...现在解决了
  • 是的,有时会发生:)

标签: java apache-kafka


【解决方案1】:

不同的答案,但在我的情况下,它恰好是消费者的初始偏移量 (auto.offset.reset)。因此,设置 auto.offset.reset=earliest 解决了我的场景中的问题。这是因为我是先发布事件,然后再启动消费者。

默认情况下,消费者只消费它启动后发布的事件,因为默认情况下auto.offset.reset=latest

例如。 consumer.properties

bootstrap.servers=localhost:9092
enable.auto.commit=true
auto.commit.interval.ms=1000
session.timeout.ms=30000
auto.offset.reset=earliest
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

测试

class KafkaEventConsumerSpecs extends FunSuite {

  case class TestEvent(eventOffset: Long, hashValue: Long, created: Date, testField: String) extends BaseEvent

  test("given an event in the event-store, consumes an event") {

    EmbeddedKafka.start()

    //PRODUCE
    val event = TestEvent(0l, 0l, new Date(), "data")
    val config = new Properties() {
      {
        load(this.getClass.getResourceAsStream("/producer.properties"))
      }
    }
    val producer = new KafkaProducer[String, String](config)

    val persistedEvent = producer.send(new ProducerRecord(event.getClass.getSimpleName, event.toString))

    assert(persistedEvent.get().offset() == 0)
    assert(persistedEvent.get().checksum() != 0)

    //CONSUME
    val consumerConfig = new Properties() {
      {
        load(this.getClass.getResourceAsStream("/consumer.properties"))
        put("group.id", "consumers_testEventsGroup")
        put("client.id", "testEventConsumer")
      }
    }

    assert(consumerConfig.getProperty("group.id") == "consumers_testEventsGroup")

    val kafkaConsumer = new KafkaConsumer[String, String](consumerConfig)

    assert(kafkaConsumer.listTopics().asScala.map(_._1).toList == List("TestEvent"))

    kafkaConsumer.subscribe(Collections.singletonList("TestEvent"))

    val events = kafkaConsumer.poll(1000)
    assert(events.count() == 1)

    EmbeddedKafka.stop()
  }
}

但是如果consumer先启动后发布,consumer应该可以消费事件而不需要auto.offset.reset设置为earliest

对 kafka 0.10 的参考

https://kafka.apache.org/documentation/#consumerconfigs

【讨论】:

    【解决方案2】:

    在我们的例子中,我们通过以下步骤解决了我们的问题:

    我们发现的第一件事是 KafkaProducer 有一个名为“重试”的配置,其默认值表示“不重试”。此外,KafkaProducer 的 send 方法是异步的,无需调用 send 方法结果的 get 方法。这样,无法保证不重试就将生成的消息传递给相应的代理。因此,您必须将其增加一点,或者可以使用 KafkaProducer 的幂等或事务模式。

    第二种情况是关于Kafka和Zookeeper版本的。我们选择了 1.0.0 版本的 Kafka 和 Zookeeper 3.4.4。特别是 Kafka 1.0.0 在与 Zookeeper 的连接方面存在问题。如果 Kafka 因意外异常而失去与 Zookeeper 的连接,它会失去对尚未同步的分区的领导权。有一个关于这个问题的错误主题: https://issues.apache.org/jira/browse/KAFKA-2729 在我们在 Kafka 日志中找到了表明与上述主题相同的问题的相应日志后,我们将 Kafka 代理版本升级到 1.1.0。

    同样重要的一点要注意,小分区(比如 100 或更小)会增加生产者的吞吐量,所以如果没有足够的消费者,那么可用的消费者就会陷入延迟消息而卡在结果上的线程中(我们用分钟测量延迟,大约 10-15 分钟)。因此,您需要根据可用资源正确平衡和配置应用程序的分区大小和线程数。

    【讨论】:

      【解决方案3】:

      还有一种情况,当新的消费者被添加到同一个组 id 时,kafka 需要很长时间才能重新平衡消费者组。 检查 kafka 日志以查看启动消费者后组是否重新平衡

      【讨论】:

        猜你喜欢
        • 2017-09-23
        • 1970-01-01
        • 1970-01-01
        • 2020-08-11
        • 2018-06-04
        • 2016-11-11
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多