Queries with streaming sources must be executed with writeStream.start(); pyspark

Viewed 168

I have trouble when trying to read the messages from kafka and the following exception appear "Queries with streaming sources must be executed with writeStream.start();"

Here my code:

from dataclasses import dataclass
from pyspark.sql import SparkSession
import pyspark.sql.functions as f

@dataclass
class DeviceData:
    device: str
    temp: float
    humd: float
    pres: float
    

spark:SparkSession = SparkSession.builder \
    .master("local[1]") \
    .appName("StreamHandler") \
    .getOrCreate()

spark.sparkContext.setLogLevel("WARN")

inputDF = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "localhost:9092") \
    .option("subscribe", "weather") \
    .load()

rawDF = inputDF.selectExpr("CAST(value AS STRING)")

df_split = inputDF.select(f.split(inputDF.value, ",")) \
    .rdd.Map(lambda x: DeviceData(x[0], x[1], x[2], x[3])) \
    .toDF(schema=['device', 'temp', 'humd', 'pres'])

summaryDF = df_split.groupBy('device') \
    .agg(f.avg('temp'), f.avg('humd'), f.avg('pres'))

query = summaryDF.writeStream.format('console').outputMode('update').start()
query.awaitTermination()

0 Answers
Related