I am going to stream data through eventhub to Databricks deltatable . the source is a very simple .net C# console program which generate random json payload . connect to azure eventhub . The evenhub has 10TU .
The source generation rate is about 10 event per seconds.
There is a databricks notebook and have some code like in the following :
val customEventhubParameters =
EventHubsConf(connStr.toString())
.setMaxEventsPerTrigger(5)
val incomingStream = spark.readStream.format("eventhubs").options(customEventhubParameters.toMap).option("maxBytesPerTrigger",100000).load()
var outputstream =
incomingStream
.select(get_json_object(($"body").cast("string"), "$.DeviceID").alias("DeviceID"), get_json_object(($"body").cast("string"), "$.detail").alias("detail"),get_json_object(($"body").cast("string"), "$.salesnumber").alias("Salesnumber").cast("Int"),get_json_object(($"body").cast("string"), "$.time").alias("Time"))
and write stream
outputstream.writeStream
.format("delta")
.outputMode("append").trigger(Trigger.ProcessingTime("10 seconds"))
.option("checkpointLocation", "/temp/checkpoint")
.start("/temp/eventdata").awaitTermination()
So it continuesly running in 10s interval. and I found that it only write 5 events to delta table at a time(which I think it's far less than the generation rate of the C# program) . How can make it a larger batch to write to delta table ??
also it cause the delta table so many small files.