【问题标题】:Can produce to Kafka but cannot consume可以生产到卡夫卡但不能消费
【发布时间】:2018-08-04 05:15:02
【问题描述】:

我正在使用 Kafka JDK 客户端版本 0.10.2.1 。我能够为 Kafka 生成简单的消息以进行“心跳”测试,但我无法使用 sdk 使用来自同一主题的消息。当我进入 Kafka CLI 时,我能够使用该消息,因此我已确认该消息在那里。这是我用来从我的 Kafka 服务器中使用的函数,带有道具 - 只有在我确实确认 produce() 成功后,我才会将生成的消息传递给主题,如果需要,我可以稍后发布该函数:

private def consumeFromKafka(topic: String, expectedMessage: String): Boolean = {
    val props: Properties = initProps("consumer")
    val consumer = new KafkaConsumer[String, String](props)
    consumer.subscribe(List(topic).asJava)
    var readExpectedRecord = false
    try {
      val records = {
        val firstPollRecs = consumer.poll(MAX_POLLTIME_MS)
        // increase timeout and try again if nothing comes back the first time in case system is busy
        if (firstPollRecs.count() == 0) firstPollRecs else {
          logger.info("KafkaHeartBeat: First poll had 0 records- trying again - doubling timeout to "
            + (MAX_POLLTIME_MS * 2)/1000 + " sec.")
          consumer.poll(MAX_POLLTIME_MS * 2)
        }
      }
      records.forEach(rec => {
        if (rec.value() == expectedMessage) readExpectedRecord = true
      })
    } catch {
      case e: Throwable => //log error
    } finally {
      consumer.close()
    }
    readExpectedRecord
  }

private def initProps(propsType: String): Properties = {
    val prop = new Properties()
    prop.put("bootstrap.servers", kafkaServer + ":" + kafkaPort)

    propsType match {
      case "producer" => {
        prop.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
        prop.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
        prop.put("acks", "1")
        prop.put("producer.type", "sync")
        prop.put("retries", "3")
        prop.put("linger.ms", "5")
      }
      case "consumer" => {
        prop.put("group.id", groupId)
        prop.put("enable.auto.commit", "false")
        prop.put("auto.commit.interval.ms", "1000")
        prop.put("session.timeout.ms", "30000")
        prop.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
        prop.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer")
        prop.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest")
        // poll just once, should only be one record for the heartbeat
        prop.put("max.poll.records", "1")
      }
    }
    prop
  }

现在,当我运行代码时,它会在控制台中输出以下内容:

13:04:21 - 发现协调器 serverName:9092 (id: 2147483647

rack: null) 组 0b8947e1-eb68-4af3-ac7b-be3f7c02e76e。 13:04:23

INFO o.a.k.c.c.i.ConsumerCoordinator - 撤销之前分配的

组 0b8947e1-eb68-4af3-ac7b-be3f7c02e76e 13:04:24 的分区 []

INFO o.a.k.c.c.i.AbstractCoordinator - (重新)加入组

0b8947e1-eb68-4af3-ac7b-be3f7c02e76e 13:04:25 信息

o.a.k.c.c.i.AbstractCoordinator - 成功加入群组

0b8947e1-eb68-4af3-ac7b-be3f7c02e76e 第 1 代 13:04:26 信息

o.a.k.c.c.i.ConsumerCoordinator - 设置新分配的分区

[HeartBeat_Topic.Service_5.2018-08-03.13_04_10.377-0] 组

0b8947e1-eb68-4af3-ac7b-be3f7c02e76e 13:04:27 信息

c.p.p.l.util.KafkaHeartBeatUtil - KafkaHeartBeat:第一次投票为 0

记录 - 再次尝试 - 将超时时间加倍至 60 秒。

然后没有别的,没有抛出错误 - 所以没有记录被轮询。有谁知道是什么阻止了“消费”的发生?订阅者似乎成功了,因为我能够成功调用 listTopics 并列出分区没有问题。

【问题讨论】:

    标签: java scala apache-kafka


    【解决方案1】:

    您的代码有错误。看来你的台词:

     if (firstPollRecs.count() == 0) 
    

    应该这样说

     if (firstPollRecs.count() > 0) 
    

    否则,您将传入一个空的firstPollRecs,然后对其进行迭代,这显然不会返回任何内容。

    【讨论】:

      猜你喜欢
      • 2020-02-20
      • 1970-01-01
      • 2019-03-27
      • 1970-01-01
      • 2020-07-28
      • 2019-08-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多