This is an automated email from the ASF dual-hosted git repository.

markli pushed a commit to branch 1.7.0-prepare
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/1.7.0-prepare by this push:
     new e04120ea2 SendCallBack to Kafka Producer
e04120ea2 is described below

commit e04120ea2df3f8130671309b34b97fb5a1dfb035
Author: Mark Li <[email protected]>
AuthorDate: Fri Feb 24 22:03:26 2023 -0800

    SendCallBack to Kafka Producer
---
 .../org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java | 6 ++++++
 1 file changed, 6 insertions(+)

diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
index a37a30352..aca8ec918 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/producer/ProducerImpl.java
@@ -19,6 +19,7 @@ package org.apache.eventmesh.connector.kafka.producer;
 
 import org.apache.eventmesh.api.RequestReplyCallback;
 import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
 import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
 
 import org.apache.kafka.clients.admin.Admin;
@@ -111,8 +112,13 @@ public class ProducerImpl {
     public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
         try {
             this.producer.send(new ProducerRecord<>(cloudEvent.getSubject(), 
cloudEvent));
+            SendResult sendResult = new SendResult();
+            sendResult.setTopic(cloudEvent.getSubject());
+            sendResult.setMessageId(cloudEvent.getId());
+            sendCallback.onSuccess(sendResult);
         } catch (Exception e) {
             log.error(String.format("Send message oneway Exception, %s", 
cloudEvent), e);
+            // 
sendCallback.onException(CloudEventUtils.convertSendResult(cloudEvent));
         }
     }
 }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to