我在 Hortonworks(HDP 2.X 版本)安装上使用 Apache Kafka。遇到的错误消息意味着 Kafka 生产者无法将数据推送到段日志文件。从命令行控制台,这意味着两件事:
- 您使用的代理端口不正确
- 您在 server.properties 中的侦听器配置不起作用
如果您在通过 scala api 编写时遇到错误消息,请另外检查使用 telnet <cluster-host> <broker-port> 与 kafka 集群的连接
注意:如果您使用 scala api 创建主题,代理需要一些时间才能了解新创建的主题。因此,在创建主题后,生产者可能会立即失败并出现错误Failed to update metadata after 60000 ms.
为了解决这个问题,我做了以下检查:
我通过 Ambari 检查后的第一个区别是 Kafka 代理在 HDP 2.x 上侦听端口 6667(apache kafka 使用 9092)。
listeners=PLAINTEXT://localhost:6667
接下来,使用 ip 而不是 localhost。
我执行了netstat -na | grep 6667
tcp 0 0 192.30.1.5:6667 0.0.0.0:* LISTEN
tcp 1 0 192.30.1.5:52242 192.30.1.5:6667 CLOSE_WAIT
tcp 0 0 192.30.1.5:54454 192.30.1.5:6667 TIME_WAIT
所以,我修改了生产者调用以使用 IP 而不是 localhost:
./kafka-console-producer.sh --broker-list 192.30.1.5:6667 --topic rdl_test_2
要监控是否有新记录正在写入,请监控/kafka-logs 文件夹。
cd /kafka-logs/<topic name>/
ls -lart
-rw-r--r--. 1 kafka hadoop 0 Feb 10 07:24 00000000000000000000.log
-rw-r--r--. 1 kafka hadoop 10485756 Feb 10 07:24 00000000000000000000.timeindex
-rw-r--r--. 1 kafka hadoop 10485760 Feb 10 07:24 00000000000000000000.index
一旦生产者写入成功,segment log-file 00000000000000000000.log 就会变大。
请看下面的尺寸:
-rw-r--r--. 1 kafka hadoop 10485760 Feb 10 07:24 00000000000000000000.index
-rw-r--r--. 1 kafka hadoop **45** Feb 10 09:16 00000000000000000000.log
-rw-r--r--. 1 kafka hadoop 10485756 Feb 10 07:24 00000000000000000000.timeindex
此时,您可以运行consumer-console.sh:
./kafka-console-consumer.sh --bootstrap-server 192.30.1.5:6667 --topic rdl_test_2 --from-beginning
response is hello world
在这一步之后,如果您想通过 Scala API 生成消息,则更改 listeners 值(从 localhost 到公共 IP)并通过 Ambari 重新启动 Kafka 代理:
listeners=PLAINTEXT://192.30.1.5:6667
样品生产者如下:
package com.scalakafka.sample
import java.util.Properties
import java.util.concurrent.TimeUnit
import org.apache.kafka.clients.producer.{ProducerRecord, KafkaProducer}
import org.apache.kafka.common.serialization.{StringSerializer, StringDeserializer}
class SampleKafkaProducer {
case class KafkaProducerConfigs(brokerList: String = "192.30.1.5:6667") {
val properties = new Properties()
val batchsize :java.lang.Integer = 1
properties.put("bootstrap.servers", brokerList)
properties.put("key.serializer", classOf[StringSerializer])
properties.put("value.serializer", classOf[StringSerializer])
// properties.put("serializer.class", classOf[StringDeserializer])
properties.put("batch.size", batchsize)
// properties.put("linger.ms", 1)
// properties.put("buffer.memory", 33554432)
}
val producer = new KafkaProducer[String, String](KafkaProducerConfigs().properties)
def produce(topic: String, messages: Iterable[String]): Unit = {
messages.foreach { m =>
println(s"Sending $topic and message is $m")
val result = producer.send(new ProducerRecord(topic, m)).get()
println(s"the write status is ${result}")
}
producer.flush()
producer.close(10L, TimeUnit.MILLISECONDS)
}
}
希望这对某人有所帮助。