This repository has been archived by the owner on Oct 12, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 22
/
EventHubReceiver.scala
executable file
·66 lines (56 loc) · 2.41 KB
/
EventHubReceiver.scala
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
// Copyright (c) Microsoft. All rights reserved.
package com.microsoft.azure.iot.kafka.connect.source
import java.time.Instant
import com.microsoft.azure.eventhubs.{EventHubClient, PartitionReceiver}
import scala.collection.JavaConverters._
import scala.collection.mutable.ListBuffer
class EventHubReceiver(val connectionString: String, val receiverConsumerGroup: String, val partition: String,
var offset: Option[String], val startTime: Option[Instant]) extends DataReceiver {
private[this] var isClosing = false
private val eventHubClient = EventHubClient.createFromConnectionStringSync(connectionString)
if (eventHubClient == null) {
throw new IllegalArgumentException("Unable to create EventHubClient from the input parameters.")
}
private val eventHubReceiver: PartitionReceiver = if (startTime.isDefined) {
eventHubClient.createReceiverSync(receiverConsumerGroup, partition.toString, startTime.get)
} else {
eventHubClient.createReceiverSync(receiverConsumerGroup, partition.toString, offset.get)
}
if (this.eventHubReceiver == null) {
throw new IllegalArgumentException("Unable to create PartitionReceiver from the input parameters.")
}
override def close(): Unit = {
if (this.eventHubReceiver != null) {
this.eventHubReceiver.synchronized {
this.isClosing = true
eventHubReceiver.close().join()
}
}
}
override def receiveData(batchSize: Int): Iterable[IotMessage] = {
var iotMessages = ListBuffer.empty[IotMessage]
var curBatchSize = batchSize
var endReached = false
// Synchronize on the eventHubReceiver object, and make sure the task is not closing,
// in which case, the eventHubReceiver might be closed.
while (curBatchSize > 0 && !endReached && !this.isClosing) {
this.eventHubReceiver.synchronized {
if(!this.isClosing) {
val batch = this.eventHubReceiver.receiveSync(curBatchSize)
if (batch != null) {
val batchIterable = batch.asScala
iotMessages ++= batchIterable.map(e => {
val content = new String(e.getBytes)
val iotDeviceData = IotMessage(content, e.getSystemProperties.asScala, e.getProperties.asScala)
iotDeviceData
})
curBatchSize -= batchIterable.size
} else {
endReached = true
}
}
}
}
iotMessages
}
}