【发布时间】: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