Create a gist now

Instantly share code, notes, and snippets.

What would you like to do?
class StatusUpdateStreamingListener(
pipelineName: String,
pipelineId: Long) extends StreamingListener {
@Inject var statusTrackerClient: StatusTrackerClient = _
override def onBatchCompleted(batch: StreamingListenerBatchCompleted): Unit = {
val msg = StreamingStatusMessage.newBuilder()
.setPipelineName(pipelineName)
.setId(pipelineId)
.setBatchTime(batch.batchInfo.batchTime.milliseconds)
.build
statusTrackerClient.send(msg)
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment