This is an automated email from the ASF dual-hosted git repository.
dubeejw pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-openwhisk.git
The following commit(s) were added to refs/heads/master by this push:
new ad57246 Wrap producer.send with Future (#3459)
ad57246 is described below
commit ad57246004bf459c46f24dd6d7e3265eabd0d39e
Author: jiangpch <[email protected]>
AuthorDate: Mon Apr 2 23:52:50 2018 +0800
Wrap producer.send with Future (#3459)
---
.../whisk/connector/kafka/KafkaProducerConnector.scala | 16 ++++++++++------
1 file changed, 10 insertions(+), 6 deletions(-)
diff --git
a/common/scala/src/main/scala/whisk/connector/kafka/KafkaProducerConnector.scala
b/common/scala/src/main/scala/whisk/connector/kafka/KafkaProducerConnector.scala
index 0c511f3..125dbef 100644
---
a/common/scala/src/main/scala/whisk/connector/kafka/KafkaProducerConnector.scala
+++
b/common/scala/src/main/scala/whisk/connector/kafka/KafkaProducerConnector.scala
@@ -36,7 +36,7 @@ import whisk.core.connector.{Message, MessageProducer}
import whisk.core.entity.UUIDs
import scala.concurrent.duration._
-import scala.concurrent.{ExecutionContext, Future, Promise}
+import scala.concurrent.{blocking, ExecutionContext, Future, Promise}
import scala.util.{Failure, Success}
class KafkaProducerConnector(kafkahosts: String, id: String =
UUIDs.randomUUID().toString)(implicit logging: Logging,
@@ -54,12 +54,16 @@ class KafkaProducerConnector(kafkahosts: String, id: String
= UUIDs.randomUUID()
val record = new ProducerRecord[String, String](topic, "messages",
msg.serialize)
val produced = Promise[RecordMetadata]()
- producer.send(record, new Callback {
- override def onCompletion(metadata: RecordMetadata, exception:
Exception): Unit = {
- if (exception == null) produced.success(metadata)
- else produced.failure(exception)
+ Future {
+ blocking {
+ producer.send(record, new Callback {
+ override def onCompletion(metadata: RecordMetadata, exception:
Exception): Unit = {
+ if (exception == null) produced.success(metadata)
+ else produced.failure(exception)
+ }
+ })
}
- })
+ }
produced.future.andThen {
case Success(status) =>
--
To stop receiving notification emails like this one, please contact
[email protected].