Forum Discussion

dt3288's avatar
dt3288
New Member
3 years ago
Solved

PySpark Notebook Using Structured Streaming with Delta Table Sink - Unsupported Operation Exception

I'm encountering difficulty reproducing the PySpark notebook example using a delta table as a streaming sink in the following training module: https://learn.microsoft.com/en-us/training/modules/work-...
  • puneetvijwani's avatar
    3 years ago

    dt3288 i have used structured steaming for incremental load before and presented on one of my sessions 
    https://www.youtube.com/watch?v=bNdKX-9nXTs
    And reference notebooks are here 
    https://github.com/puneetvijwani/fabricNotebooks

    Also i have tested your code it seems working fine for reading the testdata.csv as stream as i loaded in Files (lakehouse) and used relative path however you can also try copying abfss path by right clikcing the file and copy ABFS path

    # Welcome to your new notebook
    # Type here in the cell editor to add code!
    from pyspark.sql.types import *
    from pyspark.sql.functions import *

    # Create a stream that reads JSON data from a folder
    inputPath = 'Files/testdata.csv'
    #jsonSchema = StructType([
    csvSchema = StructType([
        StructField("device", StringType(), False),
        StructField("status", StringType(), False)
    ])
    #stream_df = spark.readStream.schema(jsonSchema).option("maxFilesPerTrigger", 1).json(inputPath)
    #stream_df = spark.readStream.schema(csvSchema).option("maxFilesPerTrigger", 1).csv(inputPath)
    stream_df = spark.readStream.format("csv").schema(csvSchema).option("header",True).option("maxFilesPerTrigger",1).load(inputPath)