【问题标题】:KafKa error java.nio.channels.UnresolvedAddressExceptionKafKa 错误 java.nio.channels.UnresolvedAddressException
【发布时间】:2018-12-24 14:58:27
【问题描述】:

我进行了 Kafka 测试,并且成功了。但是当我在 IDE 上运行程序时,我得到了这个错误,不知道如何解决它。谁能帮我?非常感谢!

public final class Constants {
    public static final String REDIS_SERVER = "localhost";

    public static final String KAFKA_SERVER = "localhost";

    public static final String KAFKA_ADDR = KAFKA_SERVER + ":9092";

    public static final String KAFKA_TOPICS = "recom1";

}

val Array(brokers, topics) = Array(Constants.KAFKA_ADDR, Constants.KAFKA_TOPICS)

val sparkConf = new 
SparkConf().setMaster("local[2]").setAppName("RealtimeRecommender")

val ssc = new StreamingContext(sparkConf, Seconds(2))

val topicsSet = topics.split(",").toSet

val kafkaParams = Map[String, String]("metadata.broker.list" -> brokers,"auto.offset.reset" -> "smallest")

val messages = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topicsSet)

我认为ip和hostname可能没有映射到/etc/hosts,但是已经有127.0.0.1 localhost。谁能帮帮我?

这是错误:

Exception in thread "main" org.apache.spark.SparkException: java.nio.channels.UnresolvedAddressException
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
    at scala.util.Either.fold(Either.scala:97)
    at org.apache.spark.streaming.kafka.KafkaCluster$.checkErrors(KafkaCluster.scala:365)
    at org.apache.spark.streaming.kafka.KafkaUtils$.getFromOffsets(KafkaUtils.scala:222)
    at org.apache.spark.streaming.kafka.KafkaUtils$.createDirectStream(KafkaUtils.scala:484)
    at com.ssx.recom.realtime.RealtimeRecommender$.main(RealtimeRecommender.scala:26)
    at com.ssx.recom.realtime.RealtimeRecommender.main(RealtimeRecommender.scala)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:606)
    at com.intellij.rt.execution.application.AppMain.main(AppMain.java:140)

【问题讨论】:

    标签: java apache-spark apache-kafka


    【解决方案1】:

    我终于解决了这个问题。只需检查zookeeper,因为zookeeper保存了kafka的配置信息,并保存了kafka的主机名。所以IP不行!

    【讨论】:

    • 如果您将 Kafka 配置为使用 IP,它将起作用。注意:Kafka DStream 在 Spark 2.3+ 中已弃用
    猜你喜欢
    • 2021-10-24
    • 1970-01-01
    • 1970-01-01
    • 2016-12-15
    • 2015-04-18
    • 1970-01-01
    • 2016-03-27
    • 2018-06-09
    • 1970-01-01
    相关资源
    最近更新 更多