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

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


The following commit(s) were added to refs/heads/master by this push:
     new 7892ab36f fix issue 2949
     new b2885b20a Merge pull request #3076 from mroccyen/issue-2949
7892ab36f is described below

commit 7892ab36f8b524f285b62a976046ced520213186
Author: mroccyen <[email protected]>
AuthorDate: Thu Feb 9 00:10:35 2023 +0800

    fix issue 2949
---
 .../rabbitmq/cloudevent/RabbitmqCloudEvent.java    | 10 ++++++---
 .../connector/rabbitmq/utils/ByteArrayUtils.java   | 25 ++++++++--------------
 2 files changed, 16 insertions(+), 19 deletions(-)

diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/cloudevent/RabbitmqCloudEvent.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/cloudevent/RabbitmqCloudEvent.java
index d0c762ac5..94a87264e 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/cloudevent/RabbitmqCloudEvent.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/cloudevent/RabbitmqCloudEvent.java
@@ -17,6 +17,8 @@
 
 package org.apache.eventmesh.connector.rabbitmq.cloudevent;
 
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.utils.JsonUtils;
 import 
org.apache.eventmesh.connector.rabbitmq.exception.RabbitmqaConnectorException;
 import org.apache.eventmesh.connector.rabbitmq.utils.ByteArrayUtils;
 
@@ -31,6 +33,8 @@ import io.cloudevents.CloudEvent;
 import io.cloudevents.SpecVersion;
 import io.cloudevents.core.builder.CloudEventBuilder;
 
+import com.fasterxml.jackson.core.type.TypeReference;
+
 import lombok.Data;
 import lombok.NoArgsConstructor;
 
@@ -70,8 +74,8 @@ public class RabbitmqCloudEvent implements Serializable {
         return optionalBytes.orElseGet(() -> new byte[]{});
     }
 
-    public static RabbitmqCloudEvent getFromByteArray(byte[] body) throws 
Exception {
-        Optional<RabbitmqCloudEvent> optionalCloudEvent = 
ByteArrayUtils.bytesToObject(body);
-        return optionalCloudEvent.orElse(null);
+    public static RabbitmqCloudEvent getFromByteArray(byte[] body) {
+        return JsonUtils.deserialize(new String(body, 
Constants.DEFAULT_CHARSET), new TypeReference<RabbitmqCloudEvent>() {
+        });
     }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/utils/ByteArrayUtils.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/utils/ByteArrayUtils.java
index 5554d459f..ad4796948 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/utils/ByteArrayUtils.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/main/java/org/apache/eventmesh/connector/rabbitmq/utils/ByteArrayUtils.java
@@ -17,33 +17,26 @@
 
 package org.apache.eventmesh.connector.rabbitmq.utils;
 
-import java.io.ByteArrayInputStream;
-import java.io.ByteArrayOutputStream;
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.utils.JsonUtils;
+
 import java.io.IOException;
-import java.io.ObjectInputStream;
-import java.io.ObjectOutputStream;
 import java.util.Optional;
 
+import com.fasterxml.jackson.core.type.TypeReference;
+
 @SuppressWarnings("all")
 public class ByteArrayUtils {
 
     public static <T> Optional<byte[]> objectToBytes(T obj) throws IOException 
{
-        byte[] bytes = null;
-        ByteArrayOutputStream out = new ByteArrayOutputStream();
-        ObjectOutputStream sOut = new ObjectOutputStream(out);
-        sOut.writeObject(obj);
-        sOut.flush();
-        bytes = out.toByteArray();
-
+        String s = JsonUtils.serialize(obj);
+        byte[] bytes = s.getBytes();
         return Optional.ofNullable(bytes);
     }
 
     public static <T> Optional<T> bytesToObject(byte[] bytes) throws 
IOException, ClassNotFoundException {
-        T t = null;
-        ByteArrayInputStream in = new ByteArrayInputStream(bytes);
-        ObjectInputStream sIn = new ObjectInputStream(in);
-        t = (T) sIn.readObject();
+        T t = JsonUtils.deserialize(new String(bytes, 
Constants.DEFAULT_CHARSET), new TypeReference<T>() {
+        });
         return Optional.ofNullable(t);
-
     }
 }


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

Reply via email to