如何用PySpark及自定义Python函数处理EventHub流?遇执行异常求助
Hey there! Let's break down the two problems you're hitting with your Spark 2.3.0 + PySpark 2.2.1 setup for streaming from Azure IoT Hub/Event Hub. I've wrestled with these exact quirks before, so here's how to get past them:
1. No foreach Sink Support in Python
Yep, you're right—structured streaming's foreach sink isn't implemented for PySpark in older versions (and even some newer ones don't have full support). The workaround here is to use foreachBatch, which was introduced in Spark 2.3.0 specifically to handle this gap for Python users.
foreachBatch lets you define a custom function that runs on every micro-batch of data from your stream. Inside this function, you can do anything you'd normally do with a static DataFrame—like use your pandas/matplotlib custom functions, write to databases, or process rows individually.
Here's a quick example to get you started:
def process_batch(batch_df, batch_id): # Do your custom processing here # Example: Convert batch to pandas DataFrame for analysis pandas_df = batch_df.toPandas() # Call your custom matplotlib plotting function your_custom_plotting_func(pandas_df) # Or write the batch to another storage system batch_df.write.mode("append").parquet("/path/to/storage") # Attach the batch processor to your streaming query streaming_query = ( your_streaming_df .writeStream .foreachBatch(process_batch) .start() ) streaming_query.awaitTermination()
This approach gives you all the flexibility you would have wanted from foreach, but works within PySpark's streaming capabilities.
2. Errors When Using .rdd, .map, or .flatMap
That error message ("Queries with streaming sources must be executed wit...") is telling you that you can't use RDD-based transformations directly on a streaming DataFrame. Structured streaming treats your data as an unbounded stream, so operations that require materializing the entire dataset (like converting to RDD) aren't allowed on the top-level stream.
There are two solid ways to fix this:
- Use
foreachBatchto handle RDD operations per batch: As with the first issue, inside yourprocess_batchfunction, thebatch_dfis a static DataFrame—so you can safely convert it to an RDD, use.map,.flatMap, etc.:def process_batch(batch_df, batch_id): # Convert batch to RDD and apply transformations rdd = batch_df.rdd.map(lambda row: (row.device_id, row.temperature * 1.8 + 32)) # Convert back to DataFrame if needed transformed_df = spark.createDataFrame(rdd, ["device_id", "temp_fahrenheit"]) # Do something with the transformed data - Use structured streaming's built-in DataFrame operations: Instead of dropping down to RDDs, try using PySpark's DataFrame API (like
select,filter,withColumnwith UDFs) which are designed to work with streams. For example, if you need to transform rows, use a pandas UDF instead of.map:from pyspark.sql.functions import pandas_udf, col import pandas as pd # Define a pandas UDF to transform temperature @pandas_udf("double") def convert_to_fahrenheit(temp_c: pd.Series) -> pd.Series: return temp_c * 1.8 + 32 # Apply the UDF directly to the streaming DataFrame transformed_stream = your_streaming_df.withColumn("temp_fahrenheit", convert_to_fahrenheit(col("temperature"))) # Then write the stream as usual query = transformed_stream.writeStream.format("console").start()
This keeps your code within the structured streaming paradigm and avoids the error entirely.
内容的提问来源于stack exchange,提问作者Eduard Biceri Matei

