Streaming from Azure Eventhub to Databricks delta table

Viewed 556

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.

0 Answers
Related