Writing a basic streaming pipeline with scala and spark. First stage output stream is being processed by second stage process. First stage output stream is saved to Azure Data Lake as parquet files. Second stage watches for these parquet files to process them further. BUT: first stage parquet files are initially created empty and only then populated with the data. This causes the second stage to processes the empty files before they are populated, so the output of second stage is empty, as it never processes any data.
Is there any configuration in spark I'm missing, to avoid this? I'm assuming this is in general not the desired behavior, as it defeats the purpose of a stream processing framework, doesn't it?
Otherwise, how can I find out that the files are complete, in the sense that they are finished being populated, so I can, for example move them to the source location where the second process is watching?
Some sample code below (the first and second stages are run simultaneously, I'm writing the code next to each other here for simplicity):
// FIRST STAGE: collect and sink to a Azure Data Lake folder.
// Here's where the files are created empty and THEN populated
val eventsStream = spark
.readStream
.format("kafka")
.option("kafka.bootstrap.servers", kafkaEndpoint)
.option("kafka.group.id", kafkaGroupId)
.option("subscribe", topic)
.option("startingOffsets", "earliest")
.load()
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
eventsStream
.writeStream
.trigger(Trigger.ProcessingTime("3 seconds"))
.option("checkpointLocation", streamWriteCollectCheckpointLocation)
.format("parquet")
.start(collectSinkFolder)
// SECOND STAGE: process first stage stream output.
// Here's where the files are processed as soon as they are created (empty) and
// the actual data is missed
val processed = spark
.readStream
.schema(dbEventSchema)
.parquet(processSourceFilesPattern)
.as[DbEvent]
.map(dbzEv => {
val utcNow: ZonedDateTime = ZonedDateTime.now(ZoneId.of("UTC"))
ProcessedEvent(dbzEv.key, utcNow.toString)
})
processed
.writeStream
.trigger(Trigger.ProcessingTime("3 seconds"))
.option("checkpointLocation", streamWriteProcessCheckpointLocation)
.format("parquet")
.start(processSinkFolder)