【问题标题】:org.apache.spark.SparkException: Writing job aborted on Databricksorg.apache.spark.SparkException:写作业在 Databricks 上中止
【发布时间】:2022-01-25 17:22:19
【问题描述】:

我使用 Databricks 从 Event Hub 中提取数据并使用 Pyspark Streaming 实时处理这些数据。代码工作正常,但在这一行之后:

df.writeStream.trigger(processingTime='100 seconds').queryName("myquery")\
  .format("console").outputMode('complete').start()

我收到以下错误:

org.apache.spark.SparkException: Writing job aborted.
Caused by: java.io.InvalidClassException: org.apache.spark.eventhubs.rdd.EventHubsRDD; local class incompatible: stream classdesc

我了解到这可能是由于处理能力低,但我使用的是 Standard_F4 机器,标准集群模式启用了自动缩放。

有什么想法吗?

【问题讨论】:

  • 您是否将格式从eventhubs 更改为console
  • 您的意思是将流读取为 eventthubs,然后将其写入控制台?我是这样做的: df = spark.readStream.format("eventhubs").options(**conf).load() 然后: df.writeStream.trigger(processingTime='100 seconds').queryName("myquery ")\ .format("console").outputMode('complete').start()
  • 我还定义了模式并将其应用于 df: import pyspark.sql.functions as F df=df.select(F.from_json(F.col("body").cast( "string"), schema).alias("streaming_df"))

标签: apache-spark pyspark azure-databricks azure-eventhub


【解决方案1】:

这看起来像是一个 JAR 问题。转到 spark 中的 JAR 文件夹,检查是否有多个 azure-eventhubs-spark_XXX.XX 的 jar。我认为您已经下载了它的不同版本并将其放置在那里,您应该从您的收藏中删除任何具有该名称的 JAR。如果您的 JAR 版本与其他 JAR 不兼容,也可能会发生此错误。尝试使用 spark config 添加 spark jars。

spark = SparkSession \
            .builder \
            .appName('my-spark') \
            .config('spark.jars.packages', 'com.microsoft.azure:azure-eventhubs-spark_2.11:2.3.12') \
            .getOrCreate()

这样spark会通过maven下载JAR文件。

【讨论】:

    猜你喜欢
    • 2020-11-07
    • 2019-12-25
    • 2018-03-18
    • 2020-12-10
    • 1970-01-01
    • 2020-10-04
    • 2022-11-02
    • 2022-09-28
    • 2022-10-12
    相关资源
    最近更新 更多