【问题标题】:Spark Streaming 1.6.0 EMR using Python : ClassNotFoundException: org.apache.spark.streaming.kinesis.KinesisUtilsPythonHelper使用 Python 的 Spark Streaming 1.6.0 EMR:ClassNotFoundException:org.apache.spark.streaming.kinesis.KinesisUtilsPythonHelper
【发布时间】:2016-03-24 20:48:34
【问题描述】:

我正在 AWS 上使用 Spark 1.6.0 和 Zeppelin 0.5.6 运行一个开箱即用的 EMR 集群。我的目标是初始化一个简单的 Spark Streaming 上下文并指向一个内部 Kinesis 流,作为概念验证。但是,当我运行它时,我得到:

Py4JJavaError: An error occurred while calling o89.loadClass. : 
java.lang.ClassNotFoundException: org.apache.spark.streaming.kinesis.KinesisUtilsPythonHelper
    at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
    at java.net.URLClassLoader$1.run(URLClassLoader.java:355)
    at java.security.AccessController.doPrivileged(Native Method)
    at java.net.URLClassLoader.findClass(URLClassLoader.java:354)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:425)
    at java.lang.ClassLoader.loadClass(ClassLoader.java:358)
    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 py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:231)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:381)
    at py4j.Gateway.invoke(Gateway.java:259)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:133)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.GatewayConnection.run(GatewayConnection.java:209)
    at java.lang.Thread.run(Thread.java:745)

我正在运行的代码(通过 Zeppelin)很简单:

%pyspark
from pyspark.streaming import StreamingContext
from pyspark.streaming.kinesis import KinesisUtils, InitialPositionInStream

ssc = StreamingContext(sc, 1)

appName = '{my-app-name}'
streamName = '{my-stream-name}'
endpointUrl = '{my-endpoint}'
regionName = '{my-region}'

lines = KinesisUtils.createStream(ssc, appName, streamName, endpointUrl, regionName, InitialPositionInStream.LATEST, 2)

当我在本地遇到这个问题时,我确保从源代码构建 spark-streaming-kinesis-asl 并将这些 jars 包含在我的 spark 配置中:

spark.driver.extraClassPath /path/to/kinesis/asl/assembly/jars/*

但是,我在 EMR 上似乎无法让它工作。为了安全起见,我将其包含在以下内容中,但无济于事:

spark.driver.extraClassPath
spark.driver.extraLibraryPath
spark.executor.extraClassPath
spark.executor.extraLibraryPath

以前有人遇到过这种情况吗?当我重新启动上下文以确认这些更改正在被拾取时,我正在打印出 spark 配置。也许这也需要在从节点上完成?或者可能是另一个配置选项/键?

【问题讨论】:

    标签: apache-spark pyspark spark-streaming amazon-emr amazon-kinesis


    【解决方案1】:

    将依赖项添加到 zeppelin 上下文“z”。下面是添加 sparkcsv 包的示例

    %dep
    z.load("com.databricks:spark-csv_2.11:1.3.0")
    

    【讨论】:

    • 这成功了!对于其他寻找更多信息的人,您可以找到完整的 Zeppelin 文档here
    猜你喜欢
    • 2016-04-25
    • 1970-01-01
    • 2015-10-18
    • 1970-01-01
    • 1970-01-01
    • 2020-03-19
    • 2016-11-29
    • 2017-08-11
    • 2021-04-03
    相关资源
    最近更新 更多