Forum Discussion
Spark Job Structured Streaming Lakehouse Files
- 2 years ago
BilalBobat lets try this.. more simpler version, This works for me just Fine
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType
spark = SparkSession.builder \
.appName("Stream CSV to Delta Table") \
.getOrCreate()
userSchema = StructType().add("name", "string").add("sales", "integer")
streamingDF = spark.readStream \
.schema(userSchema) \
.option("maxFilesPerTrigger", 1) \
.csv("Files/Streaming/") # Replace with the actual path to your streaming CSV files
query = streamingDF.writeStream \
.trigger(processingTime='5 seconds') \
.outputMode("append") \
.format("delta") \
.option("checkpointLocation", "Tables/Streaming_Table_test/_checkpoint") \
.start("Tables/Streaming_Table_test") # Replace with the path where you want to save the Delta table
query.awaitTermination()
BilalBobat
Sorry for late response , can you try this code and let me know if its working for you
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType
if __name__ == "__main__":
try:
spark = SparkSession.builder.appName("MyApp").getOrCreate()
spark.sparkContext.setLogLevel("DEBUG")
# Define the schema
userSchema = StructType().add("name", "string").add("sales", "integer")
# Read the csv files into a DataFrame
query = (spark.readStream
.schema(userSchema)
.option("maxFilesPerTrigger", 1)
.csv("Files/streamingdata/streamingfiles")
.writeStream
.trigger(processingTime='5 seconds') # Added a time-based trigger
.format("delta")
.outputMode("append")
.option("checkpointLocation", "Files/_checkpoint/Struc_streaming_csv_data")
.toTable("Struc_streaming_csv_data"))
query.awaitTermination()
except Exception as e:
print(f"An error occurred: {e}")