【问题标题】:Unable to create a stream in spark-streaming using kinesis stream无法使用 kinesis 流在火花流中创建流
【发布时间】:2019-07-19 18:53:42
【问题描述】:

我是 kinesis 新手,我正在尝试使用 spark-streaming (Pyspark) 处理 kinesis 流数据并面临以下错误

下面是我的代码:我正在将 twitter 数据推送到我的 kinesis 流并尝试使用 Spark-streaming 进行处理。我尝试包含所有依赖项的 --jars,但仍然面临同样的问题。Spark 版本 -2.4.3 和 2.3.3 以及适当的 spark-streaming-kinesis-asl-assembly.jar

from pyspark.streaming.kinesis import KinesisUtils, InitialPositionInStream
from pyspark import SparkConf,SparkContext
from pyspark.sql import SparkSession
from pyspark import StorageLevel
from pyspark.streaming import StreamingContext
from pyspark.streaming.kinesis import KinesisUtils,
InitialPositionInStream


spark_session = SparkSession.builder.getOrCreate()
ssc = StreamingContext(spark_session.sparkContext, 10)
sc = spark_session.sparkContext
Kinesis_app_name = "test"
Kinesis_stream_name = "python-stream"
endpoint_url = "https://kinesis.us-east-1.amazonaws.com"
region_name = "us-east-1"

data = KinesisUtils.createStream(
    ssc, Kinesis_app_name, Kinesis_stream_name, endpoint_url,
    region_name, InitialPositionInStream.LATEST, 10, StorageLevel.MEMORY_AND_DISK_2)


data.pprint()


ssc.start()  # Start the computation
ssc.awaitTermination()

我想使用 spark-streaming 处理流,但出现以下错误:

File "C:\spark-2.3.3-bin-hadoop2.7\python\lib\pyspark.zip\pyspark\streaming\kinesis.py", line 92, in createStream
          File "C:\spark-2.3.3-bin-hadoop2.7\python\lib\py4j-0.10.7-src.zip\py4j\java_gateway.py", line 1257, in __call__
          File "C:\spark-2.3.3-bin-hadoop2.7\python\lib\py4j-0.10.7-src.zip\py4j\protocol.py", line 328, in get_return_value
        py4j.protocol.Py4JJavaError: An error occurred while calling o27.createStream.
        : java.lang.NoClassDefFoundError: com/amazonaws/services/kinesis/model/Record
                at java.lang.Class.getDeclaredMethods0(Native Method)
                at java.lang.Class.privateGetDeclaredMethods(Class.java:2701)
                at java.lang.Class.getDeclaredMethods(Class.java:1975)
                at org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:232)
                at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:159)
                at org.apache.spark.SparkContext.clean(SparkContext.scala:2299)
                at org.apache.spark.streaming.kinesis.KinesisUtils$.createStream(KinesisUtils.scala:127)
                at org.apache.spark.streaming.kinesis.KinesisUtils$.createStream(KinesisUtils.scala:554)
                at org.apache.spark.streaming.kinesis.KinesisUtilsPythonHelper.createStream(KinesisUtils.scala:616)
                at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
                at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
                at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
                at java.lang.reflect.Method.invoke(Method.java:498)
                at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
                at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
                at py4j.Gateway.invoke(Gateway.java:282)
                at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
                at py4j.commands.CallCommand.execute(CallCommand.java:79)
                at py4j.GatewayConnection.run(GatewayConnection.java:238)
                at java.lang.Thread.run(Thread.java:748)
        Caused by: java.lang.ClassNotFoundException: com.amazonaws.services.kinesis.model.Record
                at java.net.URLClassLoader.findClass(URLClassLoader.java:382)
                at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
                at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
                ... 20 more

【问题讨论】:

    标签: pyspark spark-streaming amazon-kinesis


    【解决方案1】:

    我遇到了同样的问题。结果是我只包含了 spark-streaming-kinesis-asl jar。据我所知,这个 jar 包含 kinesis sdk。我通过删除单独的 jar 并使用带有 spark-submit 参数 --packages org.apache.spark:spark-streaming-kinesis-asl_2.12:2.4.4 的依赖项的包管理器来修复它。如果您使用包管理器但不删除有问题的 jar,则该程序将无法运行。我希望这对所有将来遇到此错误的人有所帮助。

    【讨论】:

      【解决方案2】:

      请看下面的解决方案:

      from pyspark import SparkContext
      from pyspark.streaming import StreamingContext
      from pyspark.streaming.kinesis import KinesisUtils, InitialPositionInStream, StorageLevel
      
      if __name__ == "__main__":
      
          kinesisConf = {...} # I put all my credentials in here
      
          batchInterval = 2000
          kinesisCheckpointInterval = batchInterval
      
          sc = SparkContext(appName="kinesis-stream")
          ssc = StreamingContext(sc, batchInterval)
      
          data = KinesisUtils.createStream(
              ssc=ssc,
              kinesisAppName=kinesisConf['appName'],
              streamName=kinesisConf['streamName'],
              endpointUrl=kinesisConf['endpointUrl'],
              regionName=kinesisConf['regionName'],
              initialPositionInStream=InitialPositionInStream.LATEST,
              checkpointInterval=kinesisCheckpointInterval,
              storageLevel=StorageLevel.MEMORY_AND_DISK_2,
              awsAccessKeyId=kinesisConf['awsAccessKeyId'],
              awsSecretKey=kinesisConf['awsSecretKey']
          )
      
          data.pprint()
      
          ssc.start()
          ssc.awaitTermination()
      

      & 当你运行它时,这样做:

      spark-submit --master local[8] --packages org.apache.spark:spark-streaming-kinesis-asl_2.12:3.0.0-preview ./streaming.py
      

      2.12 -> 指的是scala版本 3.0.0 -> 指的是 spark 版本

      转到here 并确保为该包选择正确的参数

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2017-04-16
        • 1970-01-01
        • 1970-01-01
        • 2018-08-09
        • 2020-03-10
        • 1970-01-01
        相关资源
        最近更新 更多