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].

Reply via email to