【问题标题】:apache spark streaming kafka integration error JAVAapache spark流kafka集成错误JAVA
【发布时间】:2018-01-09 18:41:20
【问题描述】:

我正在使用 Spark Streaming 和 kafka,但出现此错误。

线程“streaming-start”中的异常 java.lang.NoSuchMethodError: scala.Predef$.ArrowAssoc(Ljava/lang/Object;)Ljava/lang/Object; 在 org.apache.spark.streaming.kafka010.DirectKafkaInputDStream$$anonfun$start$1.apply(DirectKafkaInputDStream.scala:246) 在 org.apache.spark.streaming.kafka010.DirectKafkaInputDStream$$anonfun$start$1.apply(DirectKafkaInputDStream.scala:245) 在 scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244) 在 scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:244) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:727) 在 scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1157) 在 scala.collection.IterableLike$class.foreach(Ite​​rableLike.scala:72) 在 scala.collection.AbstractIterable.foreach(Ite​​rable.scala:54) 在 scala.collection.TraversableLike$class.map(TraversableLike.scala:244) 在 scala.collection.mutable.AbstractSet.scala$collection$SetLike$$super$map(Set.scala:45) 在 scala.collection.SetLike$class.map(SetLike.scala:93) 在 scala.collection.mutable.AbstractSet.map(Set.scala:45) 在 org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.start(DirectKafkaInputDStream.scala:245) 在 org.apache.spark.streaming.DStreamGraph$$anonfun$start$5.apply(DStreamGraph.scala:49) 在 org.apache.spark.streaming.DStreamGraph$$anonfun$start$5.apply(DStreamGraph.scala:49) 在 scala.collection.parallel.mutable.ParArray$ParArrayIterator.foreach_quick(ParArray.scala:145) 在 scala.collection.parallel.mutable.ParArray$ParArrayIterator.foreach(ParArray.scala:138) 在 scala.collection.parallel.ParIterableLike$Foreach.leaf(ParIterableLike.scala:975) 在 scala.collection.parallel.Task$$anonfun$tryLeaf$1.apply$mcV$sp(Tasks.scala:54) 在 scala.collection.parallel.Task$$anonfun$tryLeaf$1.apply(Tasks.scala:53) 在 scala.collection.parallel.Task$$anonfun$tryLeaf$1.apply(Tasks.scala:53) 在 scala.collection.parallel.Task$class.tryLeaf(Tasks.scala:56) 在 scala.collection.parallel.ParIterableLike$Foreach.tryLeaf(ParIterableLike.scala:972) 在 scala.collection.parallel.AdaptiveWorkStealingTasks$WrappedTask$class.compute(Tasks.scala:165) 在 scala.collection.parallel.AdaptiveWorkStealingForkJoinTasks$WrappedTask.compute(Tasks.scala:514) 在 scala.concurrent.forkjoin.RecursiveAction.exec(RecursiveAction.java:160) 在 scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260) 在 scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339) 在 scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979) 在 scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107) 17/08/02 16:24:58 INFO StreamingContext:StreamingContext 开始

<dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-sql_2.10</artifactId>
            <version>2.1.0</version>
        </dependency>
        <dependency>
            <groupId>com.google.code.gson</groupId>
            <artifactId>gson</artifactId>
            <version>2.8.0</version>
        </dependency>
        <dependency>
            <groupId>com.googlecode.json-simple</groupId>
            <artifactId>json-simple</artifactId>
            <version>1.1.1</version>
        </dependency>
        <dependency>
            <groupId>org.mongodb</groupId>
            <artifactId>mongo-java-driver</artifactId>
            <version>3.4.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming-kafka-0-10_2.11</artifactId>
            <version>2.2.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-core_2.10</artifactId>
            <version>2.1.0</version>
        </dependency>
        <dependency>
            <groupId>org.apache.spark</groupId>
            <artifactId>spark-streaming_2.10</artifactId>
            <version>2.1.1</version>
        </dependency>
    <dependency>
        <groupId>org.apache.kafka</groupId>
        <artifactId>kafka-clients</artifactId>
        <version>0.10.0.0</version>
    </dependency>

我的代码:

     SparkConf conf = new SparkConf().setAppName("Streaming").setMaster("local");


       JavaStreamingContext streamingContext = new JavaStreamingContext(conf, new Duration(1000));
        Map<String, Object> kafkaParams = new HashMap<>();
        kafkaParams.put("bootstrap.servers", "localhost:9092");
        kafkaParams.put("key.deserializer", StringDeserializer.class);
        kafkaParams.put("value.deserializer", StringDeserializer.class);
        kafkaParams.put("group.id", "exastax");
        kafkaParams.put("auto.offset.reset", "latest");
        kafkaParams.put("enable.auto.commit", false);

        Collection<String> topics = Arrays.asList("loglar");
        JavaInputDStream<ConsumerRecord<String, String>> stream =
                KafkaUtils.createDirectStream(
                        streamingContext,
                        LocationStrategies.PreferConsistent(),
                        ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams)
                );
stream.foreachRDD(rdd -> {
        OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
        rdd.foreachPartition(consumerRecords -> {
            OffsetRange o = offsetRanges[TaskContext.get().partitionId()];
            System.out.println(
                    o.topic() + " " + o.partition() + " " + o.fromOffset() + " " + o.untilOffset());
        });
    });
        streamingContext.start();
        streamingContext.awaitTermination();
    }
}

我正在使用 Kafka_2.11-0.11.0.0 我试图搜索这个问题,但我找不到相关的 jar。请帮我解决这个问题。

【问题讨论】:

    标签: java scala apache-spark apache-kafka spark-streaming


    【解决方案1】:

    您正在混合使用 Scala 2.10 和 Scala 2.11 代码。在 Scala 2.10 中使用 Kafka 依赖项,或在 Scala 2.11 中使用 Spark。

    【讨论】:

    • 我不明白我需要在哪里组织..是代码还是maven?
    • @TalhaK。这个 maven dep:spark-streaming-kafka-0-10_2.11 必须是 spark-streaming-kafka-0-10_2.10,就像其他 spark 依赖项一样
    • 是的!它运行了!非常感谢@maasg先生。但是我不能在Kafka中写值。 stream.foreachRDD 如何打印值?
    猜你喜欢
    • 2019-09-18
    • 2020-10-14
    • 2019-04-02
    • 2017-06-30
    • 2017-11-29
    • 2020-05-12
    • 2021-01-11
    • 1970-01-01
    • 2021-01-03
    相关资源
    最近更新 更多