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]