【问题标题】:Spark Kafka streaming in spark 2.3.0 with pythonSpark Kafka 在 spark 2.3.0 中使用 python 流式传输
【发布时间】:2019-05-19 08:58:28
【问题描述】:

我最近升级到 Spark 2.3.0。我有一个现有的 spark 作业,它曾经在 spark 2.2.0 上运行。 我正面临 AbstractMethodError 的 Java 异常

我的简单代码:

from pyspark import SparkContext                                                                                  
from pyspark.streaming import StreamingContext                                                                    
from pyspark.streaming.kafka import KafkaUtils                                                                                                                                                                                       

 if __name__ == "__main__":                                                                                        
     print "Here it is!"                                                                                               
     sc = SparkContext(appName="Tester")                                                                           
     ssc = StreamingContext(sc, 1)                                                                                 

这在 Spark 2.2.0 上运行良好

使用 Spark spark 2.3.0,我收到以下异常:

ssc = StreamingContext(sc, 1) 
  File "/usr/hdp/current/spark2-client/python/lib/pyspark.zip/pyspark/streaming/context.py", line 61, in __init__
  File "/usr/hdp/current/spark2-client/python/lib/pyspark.zip/pyspark/streaming/context.py", line 65, in _initialize_context
  File "/usr/hdp/current/spark2-client/python/lib/py4j-0.10.6-src.zip/py4j/java_gateway.py", line 1428, in __call__
  File "/usr/hdp/current/spark2-client/python/lib/py4j-0.10.6-src.zip/py4j/protocol.py", line 320, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling None.org.apache.spark.streaming.api.java.JavaStreamingContext.
: java.lang.AbstractMethodError
    at org.apache.spark.util.ListenerBus$class.$init$(ListenerBus.scala:35)
    at org.apache.spark.streaming.scheduler.StreamingListenerBus.<init>(StreamingListenerBus.scala:30)
    at org.apache.spark.streaming.scheduler.JobScheduler.<init>(JobScheduler.scala:57)
    at org.apache.spark.streaming.StreamingContext.<init>(StreamingContext.scala:184)
    at org.apache.spark.streaming.StreamingContext.<init>(StreamingContext.scala:76)
    at org.apache.spark.streaming.api.java.JavaStreamingContext.<init>(JavaStreamingContext.scala:130)
    at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
    at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
    at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
    at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:247)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
    at py4j.Gateway.invoke(Gateway.java:238)
    at py4j.commands.ConstructorCommand.invokeConstructor(ConstructorCommand.java:80)
    at py4j.commands.ConstructorCommand.execute(ConstructorCommand.java:69)
    at py4j.GatewayConnection.run(GatewayConnection.java:214)
    at java.lang.Thread.run(Thread.java:745)

我正在使用spark-streaming-kafka-0-8_2.11-2.3.0.jar 用于带有-packages 选项的spark-submit 命令。 我尝试使用 spark-streaming-kafka-0-8-assembly_2.11-2.3.0.jar 以及 --package 和 --jars 选项。

Python version: 2.7.5

我在这里遵循指南:https://spark.apache.org/docs/2.3.0/streaming-kafka-0-8-integration.html

spark streaming kafka 版本 0-8 在 2.3.0 中已弃用,但根据文档,它仍然存在。

我的命令如下:

spark-submit --master spark://10.183.0.41:7077 --packages org.apache.spark:spark-streaming-kafka-0-8_2.11:2.3.0  Kafka_test.py

当然,Spark 的 scala 底层代码发生了一些变化。

有人遇到过同样的问题吗?

【问题讨论】:

标签: apache-spark pyspark apache-kafka spark-streaming spark-streaming-kafka


【解决方案1】:

https://spark.apache.org/docs/2.3.0/streaming-kafka-integration.html

spark2.3.0 已弃用对 kafka 0.8 的支持

【讨论】:

    猜你喜欢
    • 2015-12-12
    • 2016-06-21
    • 2018-06-25
    • 1970-01-01
    • 2015-08-22
    • 1970-01-01
    • 2016-11-22
    • 1970-01-01
    • 2018-05-27
    相关资源
    最近更新 更多