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

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


The following commit(s) were added to refs/heads/cloudevents by this push:
     new 3ae9234  connector support cloud event (#586)
3ae9234 is described below

commit 3ae92347ca341a34afae7d790e5e72aa34d7e2dd
Author: wangshaojie4039 <[email protected]>
AuthorDate: Tue Nov 16 09:59:11 2021 +0800

    connector support cloud event (#586)
    
    Co-authored-by: wangshaojie <[email protected]>
---
 .../rocketmq/MessagingAccessPointImpl.java         |   7 +-
 .../cloudevent/RocketMQMessageFactory.java         |  79 +++++
 .../impl/RocketMQBinaryMessageReader.java          |  67 ++++
 .../rocketmq/cloudevent/impl/RocketMQHeaders.java  |  39 +++
 .../cloudevent/impl/RocketMQMessageWriter.java     | 104 ++++++
 .../rocketmq/consumer/PushConsumerImpl.java        | 373 +++++++++------------
 .../rocketmq/consumer/RocketMQConsumerImpl.java    |  76 +----
 .../exception/RMQMessageFormatException.java       |  18 +
 .../rocketmq/exception/RMQTimeoutException.java    |  18 +
 ...tractOMSProducer.java => AbstractProducer.java} |  59 ++--
 .../connector/rocketmq/producer/ProducerImpl.java  | 123 +++----
 .../rocketmq/producer/RocketMQProducerImpl.java    |  59 ++--
 .../connector/rocketmq/utils/CloudEventUtils.java  | 125 +++++++
 .../rocketmq/consumer/PushConsumerImplTest.java    |  41 ++-
 .../apache/rocketmq/producer/ProducerImplTest.java |  53 +--
 15 files changed, 800 insertions(+), 441 deletions(-)

diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/MessagingAccessPointImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/MessagingAccessPointImpl.java
index f1b7ec7..8770b7d 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/MessagingAccessPointImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/MessagingAccessPointImpl.java
@@ -29,9 +29,6 @@ import io.openmessaging.api.order.OrderProducer;
 import io.openmessaging.api.transaction.LocalTransactionChecker;
 import io.openmessaging.api.transaction.TransactionProducer;
 
-import org.apache.eventmesh.connector.rocketmq.consumer.PushConsumerImpl;
-import org.apache.eventmesh.connector.rocketmq.producer.ProducerImpl;
-
 public class MessagingAccessPointImpl implements MessagingAccessPoint {
 
     private Properties accessPointProperties;
@@ -52,7 +49,7 @@ public class MessagingAccessPointImpl implements 
MessagingAccessPoint {
 
     @Override
     public Producer createProducer(Properties properties) {
-        return new ProducerImpl(this.accessPointProperties);
+        return null;
     }
 
     @Override
@@ -72,7 +69,7 @@ public class MessagingAccessPointImpl implements 
MessagingAccessPoint {
 
     @Override
     public Consumer createConsumer(Properties properties) {
-        return new PushConsumerImpl(properties);
+        return null;
     }
 
     @Override
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/RocketMQMessageFactory.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/RocketMQMessageFactory.java
new file mode 100644
index 0000000..38c21d0
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/RocketMQMessageFactory.java
@@ -0,0 +1,79 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *  Unless required by applicable law or agreed to in writing, software
+ *  distributed under the License is distributed on an "AS IS" BASIS,
+ *  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *  See the License for the specific language governing permissions and
+ *  limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.rocketmq.cloudevent;
+
+import org.apache.rocketmq.common.message.Message;
+
+import java.util.Map;
+
+import javax.annotation.ParametersAreNonnullByDefault;
+
+import io.cloudevents.core.message.MessageReader;
+import io.cloudevents.core.message.MessageWriter;
+import io.cloudevents.core.message.impl.GenericStructuredMessageReader;
+import io.cloudevents.core.message.impl.MessageUtils;
+import io.cloudevents.lang.Nullable;
+import io.cloudevents.rw.CloudEventRWException;
+import io.cloudevents.rw.CloudEventWriter;
+
+import 
org.apache.eventmesh.connector.rocketmq.cloudevent.impl.RocketMQBinaryMessageReader;
+import org.apache.eventmesh.connector.rocketmq.cloudevent.impl.RocketMQHeaders;
+import 
org.apache.eventmesh.connector.rocketmq.cloudevent.impl.RocketMQMessageWriter;
+
+
+@ParametersAreNonnullByDefault
+public final class RocketMQMessageFactory {
+
+    private RocketMQMessageFactory() {
+        // prevent instantiation
+    }
+
+    public static MessageReader createReader(final Message message) throws 
CloudEventRWException {
+        return createReader(message.getProperties(), message.getBody());
+    }
+
+
+    public static MessageReader createReader(final Map<String, String> props,
+                                             @Nullable final byte[] body)
+        throws CloudEventRWException {
+
+        return MessageUtils.parseStructuredOrBinaryMessage(
+            () -> props.get(RocketMQHeaders.CONTENT_TYPE),
+            format -> new GenericStructuredMessageReader(format, body),
+            () -> props.get(RocketMQHeaders.SPEC_VERSION),
+            sv -> new RocketMQBinaryMessageReader(sv, props, body)
+        );
+    }
+
+
+    public static MessageWriter<CloudEventWriter<Message>, Message> 
createWriter(String topic) {
+        return new RocketMQMessageWriter<>(topic);
+    }
+
+    public static MessageWriter<CloudEventWriter<Message>, Message> 
createWriter(String topic,
+                                                                               
  String keys) {
+        return new RocketMQMessageWriter<>(topic, keys);
+    }
+
+    public static MessageWriter<CloudEventWriter<Message>, Message> 
createWriter(String topic,
+                                                                               
  String keys,
+                                                                               
  String tags) {
+        return new RocketMQMessageWriter<>(topic, keys, tags);
+    }
+
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQBinaryMessageReader.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQBinaryMessageReader.java
new file mode 100644
index 0000000..a02d074
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQBinaryMessageReader.java
@@ -0,0 +1,67 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *  Unless required by applicable law or agreed to in writing, software
+ *  distributed under the License is distributed on an "AS IS" BASIS,
+ *  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *  See the License for the specific language governing permissions and
+ *  limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.rocketmq.cloudevent.impl;
+
+import java.util.Map;
+import java.util.Objects;
+import java.util.function.BiConsumer;
+
+import io.cloudevents.SpecVersion;
+import io.cloudevents.core.data.BytesCloudEventData;
+import io.cloudevents.core.message.impl.BaseGenericBinaryMessageReaderImpl;
+
+public class RocketMQBinaryMessageReader
+    extends BaseGenericBinaryMessageReaderImpl<String, String> {
+
+    private final Map<String, String> headers;
+
+    public RocketMQBinaryMessageReader(SpecVersion version, Map<String, 
String> headers,
+                                       byte[] payload) {
+        super(version,
+            payload != null && payload.length > 0 ? 
BytesCloudEventData.wrap(payload) : null);
+
+        Objects.requireNonNull(headers);
+        this.headers = headers;
+    }
+
+    @Override
+    protected boolean isContentTypeHeader(String key) {
+        return key.equals(RocketMQHeaders.CONTENT_TYPE);
+    }
+
+    @Override
+    protected boolean isCloudEventsHeader(String key) {
+        return key.length() > 3 && key.substring(0, 
RocketMQHeaders.CE_PREFIX.length())
+            .startsWith(RocketMQHeaders.CE_PREFIX);
+    }
+
+    @Override
+    protected String toCloudEventsKey(String key) {
+        return key.substring(RocketMQHeaders.CE_PREFIX.length()).toLowerCase();
+    }
+
+    @Override
+    protected void forEachHeader(BiConsumer<String, String> fn) {
+        this.headers.forEach((k, v) -> fn.accept(k, v));
+    }
+
+    @Override
+    protected String toCloudEventsValue(String value) {
+        return value;
+    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQHeaders.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQHeaders.java
new file mode 100644
index 0000000..99f5edb
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQHeaders.java
@@ -0,0 +1,39 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *  Unless required by applicable law or agreed to in writing, software
+ *  distributed under the License is distributed on an "AS IS" BASIS,
+ *  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *  See the License for the specific language governing permissions and
+ *  limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.rocketmq.cloudevent.impl;
+
+import java.util.Map;
+
+import io.cloudevents.core.message.impl.MessageUtils;
+import io.cloudevents.core.v1.CloudEventV1;
+
+public class RocketMQHeaders {
+
+    public static final String CE_PREFIX = "CE_";
+
+    protected static final Map<String, String> ATTRIBUTES_TO_HEADERS =
+        MessageUtils.generateAttributesToHeadersMapping(v -> CE_PREFIX + v);
+
+    public static final String CONTENT_TYPE =
+        ATTRIBUTES_TO_HEADERS.get(CloudEventV1.DATACONTENTTYPE);
+
+    public static final String SPEC_VERSION = 
ATTRIBUTES_TO_HEADERS.get(CloudEventV1.SPECVERSION);
+
+
+}
+
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQMessageWriter.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQMessageWriter.java
new file mode 100644
index 0000000..d253069
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/cloudevent/impl/RocketMQMessageWriter.java
@@ -0,0 +1,104 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ *  Unless required by applicable law or agreed to in writing, software
+ *  distributed under the License is distributed on an "AS IS" BASIS,
+ *  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ *  See the License for the specific language governing permissions and
+ *  limitations under the License.
+ */
+
+package org.apache.eventmesh.connector.rocketmq.cloudevent.impl;
+
+import org.apache.rocketmq.common.message.Message;
+
+import io.cloudevents.CloudEventData;
+import io.cloudevents.SpecVersion;
+import io.cloudevents.core.format.EventFormat;
+import io.cloudevents.core.message.MessageWriter;
+import io.cloudevents.rw.CloudEventContextWriter;
+import io.cloudevents.rw.CloudEventRWException;
+import io.cloudevents.rw.CloudEventWriter;
+
+
+public final class RocketMQMessageWriter<R>
+    implements MessageWriter<CloudEventWriter<Message>, Message>, 
CloudEventWriter<Message> {
+
+    private Message message;
+
+
+    public RocketMQMessageWriter(String topic) {
+        message = new Message();
+        message.setTopic(topic);
+    }
+
+    public RocketMQMessageWriter(String topic, String keys) {
+        message = new Message();
+
+        message.setTopic(topic);
+
+        if (keys != null && keys.length() > 0) {
+            message.setKeys(keys);
+        }
+    }
+
+    public RocketMQMessageWriter(String topic, String keys, String tags) {
+        message = new Message();
+
+        message.setTopic(topic);
+
+        if (tags != null && tags.length() > 0) {
+            message.setTags(tags);
+        }
+
+        if (keys != null && keys.length() > 0) {
+            message.setKeys(keys);
+        }
+    }
+
+
+    @Override
+    public CloudEventContextWriter withContextAttribute(String name, String 
value)
+        throws CloudEventRWException {
+
+        String propName = RocketMQHeaders.ATTRIBUTES_TO_HEADERS.get(name);
+        if (propName == null) {
+            propName = RocketMQHeaders.CE_PREFIX + name;
+        }
+        message.putUserProperty(propName, value);
+        return this;
+    }
+
+    @Override
+    public RocketMQMessageWriter<R> create(final SpecVersion version) {
+        message.putUserProperty(RocketMQHeaders.SPEC_VERSION, 
version.toString());
+        return this;
+    }
+
+    @Override
+    public Message setEvent(final EventFormat format, final byte[] value)
+        throws CloudEventRWException {
+        message.putUserProperty(RocketMQHeaders.CONTENT_TYPE, 
format.serializedContentType());
+        message.setBody(value);
+        return message;
+    }
+
+    @Override
+    public Message end(final CloudEventData data) throws CloudEventRWException 
{
+        message.setBody(data.toBytes());
+        return message;
+    }
+
+    @Override
+    public Message end() {
+        message.setBody(null);
+        return message;
+    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/PushConsumerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/PushConsumerImpl.java
index 47565e3..7258131 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/PushConsumerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/PushConsumerImpl.java
@@ -17,21 +17,13 @@
 
 package org.apache.eventmesh.connector.rocketmq.consumer;
 
-import java.util.Map;
-import java.util.Properties;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.atomic.AtomicBoolean;
-import io.openmessaging.api.AsyncGenericMessageListener;
-import io.openmessaging.api.AsyncMessageListener;
-import io.openmessaging.api.Consumer;
-import io.openmessaging.api.GenericMessageListener;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageListener;
-import io.openmessaging.api.MessageSelector;
-import io.openmessaging.api.exception.OMSRuntimeException;
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.AsyncConsumeContext;
+import org.apache.eventmesh.api.EventListener;
 import org.apache.eventmesh.api.EventMeshAction;
-import org.apache.eventmesh.api.EventMeshAsyncConsumeContext;
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
 import org.apache.eventmesh.common.Constants;
+import 
org.apache.eventmesh.connector.rocketmq.cloudevent.RocketMQMessageFactory;
 import org.apache.eventmesh.connector.rocketmq.common.EventMeshConstants;
 import org.apache.eventmesh.connector.rocketmq.config.ClientConfig;
 import org.apache.eventmesh.connector.rocketmq.domain.NonStandardKeys;
@@ -40,17 +32,31 @@ import 
org.apache.eventmesh.connector.rocketmq.patch.EventMeshConsumeConcurrentl
 import 
org.apache.eventmesh.connector.rocketmq.patch.EventMeshMessageListenerConcurrently;
 import org.apache.eventmesh.connector.rocketmq.utils.BeanUtils;
 import org.apache.eventmesh.connector.rocketmq.utils.OMSUtil;
+import org.apache.eventmesh.connector.rocketmq.utils.CloudEventUtils;
+
 import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
 import org.apache.rocketmq.client.exception.MQClientException;
+import 
org.apache.rocketmq.client.impl.consumer.ConsumeMessageConcurrentlyService;
+import org.apache.rocketmq.client.impl.consumer.ConsumeMessageService;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.protocol.heartbeat.MessageModel;
 import org.apache.rocketmq.remoting.protocol.LanguageCode;
 
-public class PushConsumerImpl implements Consumer {
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import io.cloudevents.CloudEvent;
+import io.openmessaging.api.exception.OMSRuntimeException;
+
+public class PushConsumerImpl {
     private final DefaultMQPushConsumer rocketmqPushConsumer;
     private final Properties properties;
     private AtomicBoolean started = new AtomicBoolean(false);
-    private final Map<String, AsyncMessageListener> subscribeTable = new 
ConcurrentHashMap<>();
+    private final Map<String, EventListener> subscribeTable = new 
ConcurrentHashMap<>();
     private final ClientConfig clientConfig;
 
     public PushConsumerImpl(final Properties properties) {
@@ -58,25 +64,23 @@ public class PushConsumerImpl implements Consumer {
         this.properties = properties;
         this.clientConfig = BeanUtils.populate(properties, ClientConfig.class);
 
-//        if 
("true".equalsIgnoreCase(System.getenv("OMS_RMQ_DIRECT_NAME_SRV"))) {
-//
-//
-//        }
         String accessPoints = clientConfig.getAccessPoints();
         if (accessPoints == null || accessPoints.isEmpty()) {
-            throw new OMSRuntimeException(-1, "OMS AccessPoints is null or 
empty.");
+            throw new ConnectorRuntimeException("OMS AccessPoints is null or 
empty.");
         }
         this.rocketmqPushConsumer.setNamesrvAddr(accessPoints.replace(',', 
';'));
         String consumerGroup = clientConfig.getConsumerId();
         if (null == consumerGroup || consumerGroup.isEmpty()) {
-            throw new OMSRuntimeException(-1, "Consumer Group is necessary for 
RocketMQ, please set it.");
+            throw new ConnectorRuntimeException(
+                "Consumer Group is necessary for RocketMQ, please set it.");
         }
         this.rocketmqPushConsumer.setConsumerGroup(consumerGroup);
         
this.rocketmqPushConsumer.setMaxReconsumeTimes(clientConfig.getRmqMaxRedeliveryTimes());
         
this.rocketmqPushConsumer.setConsumeTimeout(clientConfig.getRmqMessageConsumeTimeout());
         
this.rocketmqPushConsumer.setConsumeThreadMax(clientConfig.getRmqMaxConsumeThreadNums());
         
this.rocketmqPushConsumer.setConsumeThreadMin(clientConfig.getRmqMinConsumeThreadNums());
-        
this.rocketmqPushConsumer.setMessageModel(MessageModel.valueOf(clientConfig.getMessageModel()));
+        this.rocketmqPushConsumer.setMessageModel(
+            MessageModel.valueOf(clientConfig.getMessageModel()));
 
         String consumerId = OMSUtil.buildInstanceName();
         //this.rocketmqPushConsumer.setInstanceName(consumerId);
@@ -85,109 +89,9 @@ public class PushConsumerImpl implements Consumer {
         this.rocketmqPushConsumer.setLanguage(LanguageCode.OMS);
 
         if 
(clientConfig.getMessageModel().equalsIgnoreCase(MessageModel.BROADCASTING.name()))
 {
-            rocketmqPushConsumer.registerMessageListener(new 
EventMeshMessageListenerConcurrently() {
-
-                @Override
-                public EventMeshConsumeConcurrentlyStatus 
handleMessage(MessageExt msg, EventMeshConsumeConcurrentlyContext context) {
-                    if (msg == null) {
-                        return 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
-                    }
-
-
-//                    if (!EventMeshUtil.isValidRMBTopic(msg.getTopic())) {
-//                        return 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
-//                    }
-
-                    
msg.putUserProperty(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP, 
String.valueOf(msg.getBornTimestamp()));
-                    
msg.putUserProperty(Constants.PROPERTY_MESSAGE_STORE_TIMESTAMP, 
String.valueOf(msg.getStoreTimestamp()));
-
-                    Message omsMsg = OMSUtil.msgConvert(msg);
-
-                    AsyncMessageListener listener = 
PushConsumerImpl.this.subscribeTable.get(msg.getTopic());
-
-                    if (listener == null) {
-                        throw new OMSRuntimeException(-1,
-                                String.format("The topic/queue %s isn't 
attached to this consumer", msg.getTopic()));
-                    }
-
-                    final Properties contextProperties = new Properties();
-                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
-                    EventMeshAsyncConsumeContext omsContext = new 
EventMeshAsyncConsumeContext() {
-                        @Override
-                        public void commit(EventMeshAction action) {
-                            switch (action){
-                                case CommitMessage:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS.name());
-                                    break;
-                                case ReconsumeLater:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
-                                    break;
-                                case ManualAck:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.CONSUME_FINISH.name());
-                                    break;
-                                default:
-                                    break;
-                            }
-                        }
-                    };
-                    omsContext.setAbstractContext(context);
-                    listener.consume(omsMsg, omsContext);
-
-                    return 
EventMeshConsumeConcurrentlyStatus.valueOf(contextProperties.getProperty(NonStandardKeys.MESSAGE_CONSUME_STATUS));
-                }
-            });
+            rocketmqPushConsumer.registerMessageListener(new 
BroadCastingMessageListener());
         } else {
-            rocketmqPushConsumer.registerMessageListener(new 
EventMeshMessageListenerConcurrently() {
-
-                @Override
-                public EventMeshConsumeConcurrentlyStatus 
handleMessage(MessageExt msg, EventMeshConsumeConcurrentlyContext context) {
-                    if (msg == null) {
-                        return 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
-                    }
-//                    if (!EventMeshUtil.isValidRMBTopic(msg.getTopic())) {
-//                        return 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
-//                    }
-
-                    
msg.putUserProperty(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP, 
String.valueOf(msg.getBornTimestamp()));
-                    msg.putUserProperty(EventMeshConstants.STORE_TIMESTAMP, 
String.valueOf(msg.getStoreTimestamp()));
-
-                    Message omsMsg = OMSUtil.msgConvert(msg);
-
-                    AsyncMessageListener listener = 
PushConsumerImpl.this.subscribeTable.get(msg.getTopic());
-
-                    if (listener == null) {
-                        throw new OMSRuntimeException(-1,
-                                String.format("The topic/queue %s isn't 
attached to this consumer", msg.getTopic()));
-                    }
-
-                    final Properties contextProperties = new Properties();
-
-                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
-
-                    EventMeshAsyncConsumeContext omsContext = new 
EventMeshAsyncConsumeContext() {
-                        @Override
-                        public void commit(EventMeshAction action) {
-                            switch (action) {
-                                case CommitMessage:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS.name());
-                                    break;
-                                case ReconsumeLater:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
-                                    break;
-                                case ManualAck:
-                                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
EventMeshConsumeConcurrentlyStatus.CONSUME_FINISH.name());
-                                    break;
-                                default:
-                                    break;
-                            }
-                        }
-                    };
-                    omsContext.setAbstractContext(context);
-                    listener.consume(omsMsg, omsContext);
-
-                    return 
EventMeshConsumeConcurrentlyStatus.valueOf(contextProperties.getProperty(NonStandardKeys.MESSAGE_CONSUME_STATUS));
-                }
-            });
+            rocketmqPushConsumer.registerMessageListener(new 
ClusteringMessageListener());
         }
     }
 
@@ -195,7 +99,7 @@ public class PushConsumerImpl implements Consumer {
         return properties;
     }
 
-    @Override
+
     public void start() {
         if (this.started.compareAndSet(false, true)) {
             try {
@@ -206,19 +110,19 @@ public class PushConsumerImpl implements Consumer {
         }
     }
 
-    @Override
+
     public synchronized void shutdown() {
         if (this.started.compareAndSet(true, false)) {
             this.rocketmqPushConsumer.shutdown();
         }
     }
 
-    @Override
+
     public boolean isStarted() {
         return this.started.get();
     }
 
-    @Override
+
     public boolean isClosed() {
         return !this.isStarted();
     }
@@ -227,109 +131,158 @@ public class PushConsumerImpl implements Consumer {
         return rocketmqPushConsumer;
     }
 
-//    class MessageListenerImpl implements MessageListenerConcurrently {
-//
-//        @Override
-//        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> 
rmqMsgList,
-//            ConsumeConcurrentlyContext contextRMQ) {
-//            MessageExt rmqMsg = rmqMsgList.get(0);
-//            BytesMessage omsMsg = OMSUtil.msgConvert(rmqMsg);
-//
-//            MessageListener listener = 
PushConsumerImpl.this.subscribeTable.get(rmqMsg.getTopic());
-//
-//            if (listener == null) {
-//                throw new OMSRuntimeException("-1",
-//                    String.format("The topic/queue %s isn't attached to this 
consumer", rmqMsg.getTopic()));
-//            }
-//
-//            final KeyValue contextProperties = OMS.newKeyValue();
-//            final CountDownLatch sync = new CountDownLatch(1);
-//
-//            contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS, 
ConsumeConcurrentlyStatus.RECONSUME_LATER.name());
-//
-//            MessageListener.Context context = new MessageListener.Context() {
-//                @Override
-//                public KeyValue attributes() {
-//                    return contextProperties;
-//                }
-//
-//                @Override
-//                public void ack() {
-//                    sync.countDown();
-//                    
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
-//                        ConsumeConcurrentlyStatus.CONSUME_SUCCESS.name());
-//                }
-//            };
-//            long begin = System.currentTimeMillis();
-//            listener.onReceived(omsMsg, context);
-//            long costs = System.currentTimeMillis() - begin;
-//            long timeoutMills = clientConfig.getRmqMessageConsumeTimeout() * 
60 * 1000;
-//            try {
-//                sync.await(Math.max(0, timeoutMills - costs), 
TimeUnit.MILLISECONDS);
-//            } catch (InterruptedException ignore) {
-//            }
-//
-//            return 
ConsumeConcurrentlyStatus.valueOf(contextProperties.getString(NonStandardKeys.MESSAGE_CONSUME_STATUS));
-//        }
-//    }
-
-    @Override
-    public void subscribe(String topic, String subExpression, MessageListener 
listener) {
 
+    public void subscribe(String topic, String subExpression, EventListener 
listener) {
+        this.subscribeTable.put(topic, listener);
+        try {
+            this.rocketmqPushConsumer.subscribe(topic, subExpression);
+        } catch (MQClientException e) {
+            throw new OMSRuntimeException(-1,
+                String.format("RocketMQ push consumer can't attach to %s.", 
topic));
+        }
     }
 
-    @Override
-    public void subscribe(String topic, MessageSelector selector, 
MessageListener listener) {
 
+    public void unsubscribe(String topic) {
+        this.subscribeTable.remove(topic);
+        try {
+            this.rocketmqPushConsumer.unsubscribe(topic);
+        } catch (Exception e) {
+            throw new OMSRuntimeException(-1,
+                String.format("RocketMQ push consumer fails to unsubscribe 
topic: %s", topic));
+        }
     }
 
-    @Override
-    public <T> void subscribe(String topic, String subExpression, 
GenericMessageListener<T> listener) {
-
+    public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext 
context) {
+        ConsumeMessageService consumeMessageService = rocketmqPushConsumer
+            .getDefaultMQPushConsumerImpl().getConsumeMessageService();
+        List<MessageExt> msgExtList = new ArrayList<>(cloudEvents.size());
+        for (CloudEvent msg : cloudEvents) {
+            msgExtList.add(CloudEventUtils.msgConvertExt(
+                
RocketMQMessageFactory.createWriter(msg.getSubject()).writeBinary(msg)));
+        }
+        ((ConsumeMessageConcurrentlyService) consumeMessageService)
+            .updateOffset(msgExtList, (EventMeshConsumeConcurrentlyContext) 
context);
     }
 
-    @Override
-    public <T> void subscribe(String topic, MessageSelector selector, 
GenericMessageListener<T> listener) {
 
-    }
+    private class BroadCastingMessageListener extends 
EventMeshMessageListenerConcurrently {
 
-    @Override
-    public void subscribe(String topic, String subExpression, 
AsyncMessageListener listener) {
-        this.subscribeTable.put(topic, listener);
-        try {
-            this.rocketmqPushConsumer.subscribe(topic, subExpression);
-        } catch (MQClientException e) {
-            throw new OMSRuntimeException(-1, String.format("RocketMQ push 
consumer can't attach to %s.", topic));
+        @Override
+        public EventMeshConsumeConcurrentlyStatus handleMessage(MessageExt msg,
+                                                                
EventMeshConsumeConcurrentlyContext context) {
+            if (msg == null) {
+                return EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
+            }
+
+            msg.putUserProperty(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP,
+                String.valueOf(msg.getBornTimestamp()));
+            msg.putUserProperty(Constants.PROPERTY_MESSAGE_STORE_TIMESTAMP,
+                String.valueOf(msg.getStoreTimestamp()));
+
+            CloudEvent cloudEvent =
+                
RocketMQMessageFactory.createReader(CloudEventUtils.msgConvert(msg)).toEvent();
+
+            EventListener listener = 
PushConsumerImpl.this.subscribeTable.get(msg.getTopic());
+
+            if (listener == null) {
+                throw new OMSRuntimeException(-1,
+                    String.format("The topic/queue %s isn't attached to this 
consumer",
+                        msg.getTopic()));
+            }
+
+            final Properties contextProperties = new Properties();
+            contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
+            AsyncConsumeContext asyncConsumeContext = new 
AsyncConsumeContext() {
+                @Override
+                public void commit(EventMeshAction action) {
+                    switch (action) {
+                        case CommitMessage:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS.name());
+                            break;
+                        case ReconsumeLater:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
+                            break;
+                        case ManualAck:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.CONSUME_FINISH.name());
+                            break;
+                        default:
+                            break;
+                    }
+                }
+            };
+
+            listener.consume(cloudEvent, asyncConsumeContext);
+
+            return EventMeshConsumeConcurrentlyStatus.valueOf(
+                
contextProperties.getProperty(NonStandardKeys.MESSAGE_CONSUME_STATUS));
         }
-    }
 
-    @Override
-    public void subscribe(String topic, MessageSelector selector, 
AsyncMessageListener listener) {
 
     }
 
-    @Override
-    public <T> void subscribe(String topic, String subExpression, 
AsyncGenericMessageListener<T> listener) {
+    private class ClusteringMessageListener extends 
EventMeshMessageListenerConcurrently {
 
-    }
+        @Override
+        public EventMeshConsumeConcurrentlyStatus handleMessage(MessageExt msg,
+                                                                
EventMeshConsumeConcurrentlyContext context) {
+            if (msg == null) {
+                return EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS;
+            }
 
-    @Override
-    public <T> void subscribe(String topic, MessageSelector selector, 
AsyncGenericMessageListener<T> listener) {
+            msg.putUserProperty(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP,
+                String.valueOf(msg.getBornTimestamp()));
+            msg.putUserProperty(EventMeshConstants.STORE_TIMESTAMP,
+                String.valueOf(msg.getStoreTimestamp()));
 
-    }
+            CloudEvent cloudEvent =
+                
RocketMQMessageFactory.createReader(CloudEventUtils.msgConvert(msg)).toEvent();
 
-    @Override
-    public void unsubscribe(String topic) {
-        this.subscribeTable.remove(topic);
-        try {
-            this.rocketmqPushConsumer.unsubscribe(topic);
-        } catch (Exception e) {
-            throw new OMSRuntimeException(-1, String.format("RocketMQ push 
consumer fails to unsubscribe topic: %s", topic));
+            EventListener listener = 
PushConsumerImpl.this.subscribeTable.get(msg.getTopic());
+
+            if (listener == null) {
+                throw new OMSRuntimeException(-1,
+                    String.format("The topic/queue %s isn't attached to this 
consumer",
+                        msg.getTopic()));
+            }
+
+            final Properties contextProperties = new Properties();
+
+            contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
+
+            AsyncConsumeContext asyncConsumeContext = new 
AsyncConsumeContext() {
+                @Override
+                public void commit(EventMeshAction action) {
+                    switch (action) {
+                        case CommitMessage:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.CONSUME_SUCCESS.name());
+                            break;
+                        case ReconsumeLater:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.RECONSUME_LATER.name());
+                            break;
+                        case ManualAck:
+                            
contextProperties.put(NonStandardKeys.MESSAGE_CONSUME_STATUS,
+                                
EventMeshConsumeConcurrentlyStatus.CONSUME_FINISH.name());
+                            break;
+                        default:
+                            break;
+                    }
+                }
+            };
+
+            listener.consume(cloudEvent, asyncConsumeContext);
+
+            return EventMeshConsumeConcurrentlyStatus.valueOf(
+                
contextProperties.getProperty(NonStandardKeys.MESSAGE_CONSUME_STATUS));
         }
     }
 
-    @Override
-    public void updateCredential(Properties credentialProperties) {
 
-    }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/RocketMQConsumerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/RocketMQConsumerImpl.java
index d103cea..9626e2c 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/RocketMQConsumerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/consumer/RocketMQConsumerImpl.java
@@ -18,6 +18,8 @@
 package org.apache.eventmesh.connector.rocketmq.consumer;
 
 import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.consumer.Consumer;
 import org.apache.eventmesh.api.consumer.MeshMQPushConsumer;
 import org.apache.eventmesh.connector.rocketmq.MessagingAccessPointImpl;
 import org.apache.eventmesh.connector.rocketmq.common.Constants;
@@ -40,6 +42,7 @@ import java.util.Properties;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import io.cloudevents.CloudEvent;
 import io.openmessaging.api.AsyncGenericMessageListener;
 import io.openmessaging.api.AsyncMessageListener;
 import io.openmessaging.api.GenericMessageListener;
@@ -48,7 +51,7 @@ import io.openmessaging.api.MessageListener;
 import io.openmessaging.api.MessageSelector;
 import io.openmessaging.api.MessagingAccessPoint;
 
-public class RocketMQConsumerImpl implements MeshMQPushConsumer {
+public class RocketMQConsumerImpl implements Consumer {
 
     public Logger logger = LoggerFactory.getLogger(this.getClass());
 
@@ -72,9 +75,9 @@ public class RocketMQConsumerImpl implements 
MeshMQPushConsumer {
             consumerGroup = Constants.BROADCAST_PREFIX + consumerGroup;
         }
 
-        String omsNamesrv = clientConfiguration.namesrvAddr;
+        String namesrvAddr = clientConfiguration.namesrvAddr;
         Properties properties = new Properties();
-        properties.put("ACCESS_POINTS", omsNamesrv);
+        properties.put("ACCESS_POINTS", namesrvAddr);
         properties.put("REGION", "namespace");
         properties.put("instanceName", instanceName);
         properties.put("CONSUMER_ID", consumerGroup);
@@ -83,12 +86,12 @@ public class RocketMQConsumerImpl implements 
MeshMQPushConsumer {
         } else {
             properties.put("MESSAGE_MODEL", MessageModel.CLUSTERING.name());
         }
-        MessagingAccessPoint messagingAccessPoint = new 
MessagingAccessPointImpl(properties);
-        pushConsumer = (PushConsumerImpl) 
messagingAccessPoint.createConsumer(properties);
+
+        pushConsumer = new PushConsumerImpl(properties);
     }
 
     @Override
-    public void subscribe(String topic, AsyncMessageListener listener) throws 
Exception {
+    public void subscribe(String topic, EventListener listener) throws 
Exception {
         pushConsumer.subscribe(topic, "*", listener);
     }
 
@@ -108,16 +111,8 @@ public class RocketMQConsumerImpl implements 
MeshMQPushConsumer {
     }
 
     @Override
-    public void updateOffset(List<Message> msgs, AbstractContext context) {
-        ConsumeMessageService consumeMessageService =
-            
pushConsumer.getRocketmqPushConsumer().getDefaultMQPushConsumerImpl()
-                .getConsumeMessageService();
-        List<MessageExt> msgExtList = new ArrayList<>(msgs.size());
-        for (Message msg : msgs) {
-            msgExtList.add(OMSUtil.msgConvertExt(msg));
-        }
-        ((ConsumeMessageConcurrentlyService) consumeMessageService)
-            .updateOffset(msgExtList, (EventMeshConsumeConcurrentlyContext) 
context);
+    public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext 
context) {
+        pushConsumer.updateOffset(cloudEvents, context);
     }
 
     @Override
@@ -134,55 +129,6 @@ public class RocketMQConsumerImpl implements 
MeshMQPushConsumer {
         return pushConsumer.attributes();
     }
 
-    @Override
-    public void subscribe(String topic, String subExpression, MessageListener 
listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public void subscribe(String topic, MessageSelector selector, 
MessageListener listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public <T> void subscribe(String topic, String subExpression,
-                              GenericMessageListener<T> listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public <T> void subscribe(String topic, MessageSelector selector,
-                              GenericMessageListener<T> listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public void subscribe(String topic, String subExpression, 
AsyncMessageListener listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public void subscribe(String topic, MessageSelector selector, 
AsyncMessageListener listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public <T> void subscribe(String topic, String subExpression,
-                              AsyncGenericMessageListener<T> listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public <T> void subscribe(String topic, MessageSelector selector,
-                              AsyncGenericMessageListener<T> listener) {
-        throw new UnsupportedOperationException("not supported yet");
-    }
-
-    @Override
-    public void updateCredential(Properties credentialProperties) {
-
-    }
-
     private String getRocketMqConfigFile() {
         // get from classpath
         String configFile = RocketMQConsumerImpl.class.getClassLoader()
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQMessageFormatException.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQMessageFormatException.java
new file mode 100644
index 0000000..171f8ce
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQMessageFormatException.java
@@ -0,0 +1,18 @@
+package org.apache.eventmesh.connector.rocketmq.exception;
+
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+
+public class RMQMessageFormatException extends ConnectorRuntimeException {
+
+    public RMQMessageFormatException(String message) {
+        super(message);
+    }
+
+    public RMQMessageFormatException(Throwable throwable) {
+        super(throwable);
+    }
+
+    public RMQMessageFormatException(String message, Throwable throwable) {
+        super(message, throwable);
+    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQTimeoutException.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQTimeoutException.java
new file mode 100644
index 0000000..3d0169e
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/exception/RMQTimeoutException.java
@@ -0,0 +1,18 @@
+package org.apache.eventmesh.connector.rocketmq.exception;
+
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+
+public class RMQTimeoutException extends ConnectorRuntimeException {
+
+    public RMQTimeoutException(String message) {
+        super(message);
+    }
+
+    public RMQTimeoutException(Throwable throwable) {
+        super(throwable);
+    }
+
+    public RMQTimeoutException(String message, Throwable throwable) {
+        super(message, throwable);
+    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractOMSProducer.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractProducer.java
similarity index 69%
rename from 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractOMSProducer.java
rename to 
eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractProducer.java
index 548077a..3b0b7ce 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractOMSProducer.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/AbstractProducer.java
@@ -20,13 +20,13 @@ package org.apache.eventmesh.connector.rocketmq.producer;
 import java.util.Properties;
 import java.util.concurrent.atomic.AtomicBoolean;
 
-import io.openmessaging.api.exception.OMSMessageFormatException;
-import io.openmessaging.api.exception.OMSRuntimeException;
-import io.openmessaging.api.exception.OMSTimeOutException;
-
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
 import org.apache.eventmesh.connector.rocketmq.config.ClientConfig;
+import 
org.apache.eventmesh.connector.rocketmq.exception.RMQMessageFormatException;
+import org.apache.eventmesh.connector.rocketmq.exception.RMQTimeoutException;
 import org.apache.eventmesh.connector.rocketmq.utils.BeanUtils;
 import org.apache.eventmesh.connector.rocketmq.utils.OMSUtil;
+
 import org.apache.rocketmq.client.exception.MQBrokerException;
 import org.apache.rocketmq.client.exception.MQClientException;
 import org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl;
@@ -38,7 +38,7 @@ import 
org.apache.rocketmq.remoting.exception.RemotingConnectException;
 import org.apache.rocketmq.remoting.exception.RemotingTimeoutException;
 import org.apache.rocketmq.remoting.protocol.LanguageCode;
 
-public abstract class AbstractOMSProducer {
+public abstract class AbstractProducer {
 
     final static InternalLogger log = ClientLogger.getLog();
     final Properties properties;
@@ -48,14 +48,14 @@ public abstract class AbstractOMSProducer {
     private final ClientConfig clientConfig;
     private final String PRODUCER_ID = "PRODUCER_ID";
 
-    AbstractOMSProducer(final Properties properties) {
+    AbstractProducer(final Properties properties) {
         this.properties = properties;
         this.rocketmqProducer = new DefaultMQProducer();
         this.clientConfig = BeanUtils.populate(properties, ClientConfig.class);
 
         String accessPoints = clientConfig.getAccessPoints();
         if (accessPoints == null || accessPoints.isEmpty()) {
-            throw new OMSRuntimeException(-1, "OMS AccessPoints is null or 
empty.");
+            throw new ConnectorRuntimeException("OMS AccessPoints is null or 
empty.");
         }
 
         this.rocketmqProducer.setNamesrvAddr(accessPoints.replace(',', ';'));
@@ -75,7 +75,7 @@ public abstract class AbstractOMSProducer {
             try {
                 this.rocketmqProducer.start();
             } catch (MQClientException e) {
-                throw new OMSRuntimeException("-1", e);
+                throw new ConnectorRuntimeException("-1", e);
             }
         }
         this.started.set(true);
@@ -96,25 +96,30 @@ public abstract class AbstractOMSProducer {
         return !this.isStarted();
     }
 
-    OMSRuntimeException checkProducerException(String topic, String msgId, 
Throwable e) {
+    ConnectorRuntimeException checkProducerException(String topic, String 
msgId, Throwable e) {
         if (e instanceof MQClientException) {
             if (e.getCause() != null) {
                 if (e.getCause() instanceof RemotingTimeoutException) {
-                    return new OMSTimeOutException(-1, String.format("Send 
message to broker timeout, %dms, Topic=%s, msgId=%s",
+                    return new RMQTimeoutException(
+                        String.format("Send message to broker timeout, %dms, 
Topic=%s, msgId=%s",
                             this.rocketmqProducer.getSendMsgTimeout(), topic, 
msgId), e);
-                } else if (e.getCause() instanceof MQBrokerException || 
e.getCause() instanceof RemotingConnectException) {
+                } else if (e.getCause() instanceof MQBrokerException ||
+                    e.getCause() instanceof RemotingConnectException) {
                     if (e.getCause() instanceof MQBrokerException) {
                         MQBrokerException brokerException = 
(MQBrokerException) e.getCause();
-                        return new OMSRuntimeException(-1, 
String.format("Received a broker exception, Topic=%s, msgId=%s, %s",
+                        return new ConnectorRuntimeException(
+                            String.format("Received a broker exception, 
Topic=%s, msgId=%s, %s",
                                 topic, msgId, 
brokerException.getErrorMessage()), e);
                     }
 
                     if (e.getCause() instanceof RemotingConnectException) {
-                        RemotingConnectException connectException = 
(RemotingConnectException) e.getCause();
-                        return new OMSRuntimeException(-1,
-                                String.format("Network connection experiences 
failures. Topic=%s, msgId=%s, %s",
-                                        topic, msgId, 
connectException.getMessage()),
-                                e);
+                        RemotingConnectException connectException =
+                            (RemotingConnectException) e.getCause();
+                        return new ConnectorRuntimeException(
+                            String.format(
+                                "Network connection experiences failures. 
Topic=%s, msgId=%s, %s",
+                                topic, msgId, connectException.getMessage()),
+                            e);
                     }
                 }
             }
@@ -122,25 +127,33 @@ public abstract class AbstractOMSProducer {
             else {
                 MQClientException clientException = (MQClientException) e;
                 if (-1 == clientException.getResponseCode()) {
-                    return new OMSRuntimeException(-1, String.format("Topic 
does not exist, Topic=%s, msgId=%s",
+                    return new ConnectorRuntimeException(
+                        String.format("Topic does not exist, Topic=%s, 
msgId=%s",
                             topic, msgId), e);
                 } else if (ResponseCode.MESSAGE_ILLEGAL == 
clientException.getResponseCode()) {
-                    return new OMSMessageFormatException(-1, String.format("A 
illegal message for RocketMQ, Topic=%s, msgId=%s",
+                    return new RMQMessageFormatException(
+                        String.format("A illegal message for RocketMQ, 
Topic=%s, msgId=%s",
                             topic, msgId), e);
                 }
             }
         }
-        return new OMSRuntimeException(-1, "Send message to RocketMQ broker 
failed.", e);
+        return new ConnectorRuntimeException("Send message to RocketMQ broker 
failed.", e);
     }
 
     protected void checkProducerServiceState(DefaultMQProducerImpl producer) {
         switch (producer.getServiceState()) {
             case CREATE_JUST:
-                throw new OMSRuntimeException(String.format("You do not have 
start the producer, %s", producer.getServiceState()));
+                throw new ConnectorRuntimeException(
+                    String.format("You do not have start the producer, %s",
+                        producer.getServiceState()));
             case SHUTDOWN_ALREADY:
-                throw new OMSRuntimeException(String.format("Your producer has 
been shut down, %s", producer.getServiceState()));
+                throw new ConnectorRuntimeException(
+                    String.format("Your producer has been shut down, %s",
+                        producer.getServiceState()));
             case START_FAILED:
-                throw new OMSRuntimeException(String.format("When you start 
your service throws an exception, %s", producer.getServiceState()));
+                throw new ConnectorRuntimeException(
+                    String.format("When you start your service throws an 
exception, %s",
+                        producer.getServiceState()));
             case RUNNING:
             default:
         }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
index dbad619..ec04b58 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/ProducerImpl.java
@@ -18,30 +18,32 @@
 package org.apache.eventmesh.connector.rocketmq.producer;
 
 import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+import org.apache.eventmesh.api.exception.OnExceptionContext;
+import 
org.apache.eventmesh.connector.rocketmq.cloudevent.RocketMQMessageFactory;
 import org.apache.eventmesh.connector.rocketmq.utils.OMSUtil;
+import org.apache.eventmesh.connector.rocketmq.utils.CloudEventUtils;
 
 import org.apache.rocketmq.client.exception.MQBrokerException;
 import org.apache.rocketmq.client.exception.MQClientException;
 import org.apache.rocketmq.client.producer.RequestCallback;
+import org.apache.rocketmq.common.MixAll;
+import org.apache.rocketmq.common.message.Message;
 import org.apache.rocketmq.common.message.MessageClientIDSetter;
+import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.remoting.exception.RemotingException;
 
 import java.util.Properties;
-import java.util.concurrent.ExecutorService;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageBuilder;
-import io.openmessaging.api.OnExceptionContext;
-import io.openmessaging.api.Producer;
-import io.openmessaging.api.SendCallback;
-import io.openmessaging.api.SendResult;
-import io.openmessaging.api.exception.OMSRuntimeException;
+import io.cloudevents.CloudEvent;
 
-public class ProducerImpl extends AbstractOMSProducer implements Producer {
+public class ProducerImpl extends AbstractProducer {
 
     public static final int eventMeshServerAsyncAccumulationThreshold = 1000;
 
@@ -66,78 +68,95 @@ public class ProducerImpl extends AbstractOMSProducer 
implements Producer {
         super.getRocketmqProducer().setCompressMsgBodyOverHowmuch(10);
     }
 
-    @Override
-    public SendResult send(Message message) {
-        
this.checkProducerServiceState(rocketmqProducer.getDefaultMQProducerImpl());
-        org.apache.rocketmq.common.message.Message msgRMQ = 
OMSUtil.msgConvert(message);
 
+    public SendResult send(CloudEvent cloudEvent) {
+        
this.checkProducerServiceState(rocketmqProducer.getDefaultMQProducerImpl());
+        org.apache.rocketmq.common.message.Message msg =
+            
RocketMQMessageFactory.createWriter(cloudEvent.getSubject()).writeBinary(cloudEvent);
+        String messageId = null;
         try {
             org.apache.rocketmq.client.producer.SendResult sendResultRmq =
-                this.rocketmqProducer.send(msgRMQ);
-            message.setMsgID(sendResultRmq.getMsgId());
+                this.rocketmqProducer.send(msg);
             SendResult sendResult = new SendResult();
             sendResult.setTopic(sendResultRmq.getMessageQueue().getTopic());
-            sendResult.setMessageId(sendResultRmq.getMsgId());
+            messageId = sendResultRmq.getMsgId();
+            sendResult.setMessageId(messageId);
             return sendResult;
         } catch (Exception e) {
-            log.error(String.format("Send message Exception, %s", message), e);
-            throw this.checkProducerException(message.getTopic(), 
message.getMsgID(), e);
+            log.error(String.format("Send message Exception, %s", msg), e);
+            throw this.checkProducerException(msg.getTopic(), messageId, e);
         }
     }
 
-    @Override
-    public void sendOneway(Message message) {
-        
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
-        org.apache.rocketmq.common.message.Message msgRMQ = 
OMSUtil.msgConvert(message);
 
+    public void sendOneway(CloudEvent cloudEvent) {
+        
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
+        org.apache.rocketmq.common.message.Message msg =
+            
RocketMQMessageFactory.createWriter(cloudEvent.getSubject()).writeBinary(cloudEvent);
         try {
-            this.rocketmqProducer.sendOneway(msgRMQ);
-            message.setMsgID(MessageClientIDSetter.getUniqID(msgRMQ));
+            this.rocketmqProducer.sendOneway(msg);
         } catch (Exception e) {
-            log.error(String.format("Send message oneway Exception, %s", 
message), e);
-            throw this.checkProducerException(message.getTopic(), 
message.getMsgID(), e);
+            log.error(String.format("Send message oneway Exception, %s", msg), 
e);
+            throw this.checkProducerException(msg.getTopic(), 
MessageClientIDSetter.getUniqID(msg),
+                e);
         }
     }
 
-    @Override
-    public void sendAsync(Message message, SendCallback sendCallback) {
+
+    public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
         
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
-        org.apache.rocketmq.common.message.Message msgRMQ = 
OMSUtil.msgConvert(message);
+        org.apache.rocketmq.common.message.Message msg =
+            
RocketMQMessageFactory.createWriter(cloudEvent.getSubject()).writeBinary(cloudEvent);
 
         try {
-            this.rocketmqProducer.send(msgRMQ, 
this.sendCallbackConvert(message, sendCallback));
-            message.setMsgID(MessageClientIDSetter.getUniqID(msgRMQ));
+            this.rocketmqProducer.send(msg, this.sendCallbackConvert(msg, 
sendCallback));
         } catch (Exception e) {
-            log.error(String.format("Send message async Exception, %s", 
message), e);
-            throw this.checkProducerException(message.getTopic(), 
message.getMsgID(), e);
+            log.error(String.format("Send message async Exception, %s", msg), 
e);
+            throw this.checkProducerException(msg.getTopic(), 
MessageClientIDSetter.getUniqID(msg),
+                e);
         }
     }
 
-    public void request(Message message, RRCallback rrCallback, long timeout)
+    public void request(CloudEvent cloudEvent, RRCallback rrCallback, long 
timeout)
         throws InterruptedException, RemotingException, MQClientException, 
MQBrokerException {
 
         
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
-        org.apache.rocketmq.common.message.Message msgRMQ = 
OMSUtil.msgConvert(message);
-        rocketmqProducer.request(msgRMQ, rrCallbackConvert(message, 
rrCallback), timeout);
+        org.apache.rocketmq.common.message.Message msg =
+            
RocketMQMessageFactory.createWriter(cloudEvent.getSubject()).writeBinary(cloudEvent);
+        rocketmqProducer.request(msg, rrCallbackConvert(msg, rrCallback), 
timeout);
+    }
+
+    public boolean reply(final CloudEvent cloudEvent, final SendCallback 
sendCallback) {
+        
this.checkProducerServiceState(this.rocketmqProducer.getDefaultMQProducerImpl());
+        org.apache.rocketmq.common.message.Message msg =
+            
RocketMQMessageFactory.createWriter(cloudEvent.getSubject()).writeBinary(cloudEvent);
+        msg.putUserProperty(MessageConst.PROPERTY_MESSAGE_TYPE, 
MixAll.REPLY_MESSAGE_FLAG);
+        try {
+            this.rocketmqProducer.send(msg, this.sendCallbackConvert(msg, 
sendCallback));
+        } catch (Exception e) {
+            log.error(String.format("Send message async Exception, %s", msg), 
e);
+            throw this.checkProducerException(msg.getTopic(), 
MessageClientIDSetter.getUniqID(msg),
+                e);
+        }
+        return true;
+
     }
 
     private RequestCallback rrCallbackConvert(final Message message, final 
RRCallback rrCallback) {
         return new RequestCallback() {
             @Override
             public void onSuccess(org.apache.rocketmq.common.message.Message 
message) {
-                Message openMessage = OMSUtil.msgConvert((MessageExt) message);
+                io.openmessaging.api.Message openMessage = 
OMSUtil.msgConvert((MessageExt) message);
                 rrCallback.onSuccess(openMessage);
             }
 
             @Override
             public void onException(Throwable e) {
                 String topic = message.getTopic();
-                String msgId = message.getMsgID();
-                OMSRuntimeException onsEx =
-                    ProducerImpl.this.checkProducerException(topic, msgId, e);
+                ConnectorRuntimeException onsEx =
+                    ProducerImpl.this.checkProducerException(topic, null, e);
                 OnExceptionContext context = new OnExceptionContext();
                 context.setTopic(topic);
-                context.setMessageId(msgId);
                 context.setException(onsEx);
                 rrCallback.onException(e);
 
@@ -151,18 +170,16 @@ public class ProducerImpl extends AbstractOMSProducer 
implements Producer {
             new org.apache.rocketmq.client.producer.SendCallback() {
                 @Override
                 public void 
onSuccess(org.apache.rocketmq.client.producer.SendResult sendResult) {
-                    
sendCallback.onSuccess(OMSUtil.sendResultConvert(sendResult));
+                    
sendCallback.onSuccess(CloudEventUtils.convertSendResult(sendResult));
                 }
 
                 @Override
                 public void onException(Throwable e) {
                     String topic = message.getTopic();
-                    String msgId = message.getMsgID();
-                    OMSRuntimeException onsEx =
-                        ProducerImpl.this.checkProducerException(topic, msgId, 
e);
+                    ConnectorRuntimeException onsEx =
+                        ProducerImpl.this.checkProducerException(topic, null, 
e);
                     OnExceptionContext context = new OnExceptionContext();
                     context.setTopic(topic);
-                    context.setMessageId(msgId);
                     context.setException(onsEx);
                     sendCallback.onException(context);
                 }
@@ -170,18 +187,4 @@ public class ProducerImpl extends AbstractOMSProducer 
implements Producer {
         return rmqSendCallback;
     }
 
-    @Override
-    public void setCallbackExecutor(ExecutorService callbackExecutor) {
-//        this.rocketmqProducer.setCallbackExecutor(callbackExecutor);
-    }
-
-    @Override
-    public void updateCredential(Properties credentialProperties) {
-
-    }
-
-    @Override
-    public <T> MessageBuilder<T> messageBuilder() {
-        return null;
-    }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
index 75e2360..52af730 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/producer/RocketMQProducerImpl.java
@@ -17,33 +17,29 @@
 
 package org.apache.eventmesh.connector.rocketmq.producer;
 
-import java.io.File;
-import java.util.Properties;
-import java.util.concurrent.ExecutorService;
-
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageBuilder;
-import io.openmessaging.api.MessagingAccessPoint;
-import io.openmessaging.api.OMS;
-import io.openmessaging.api.OMSBuiltinKeys;
-import io.openmessaging.api.SendCallback;
-import io.openmessaging.api.SendResult;
-
 import org.apache.eventmesh.api.RRCallback;
-import org.apache.eventmesh.api.producer.MeshMQProducer;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.producer.Producer;
 import org.apache.eventmesh.connector.rocketmq.MessagingAccessPointImpl;
 import org.apache.eventmesh.connector.rocketmq.common.EventMeshConstants;
 import org.apache.eventmesh.connector.rocketmq.config.ClientConfiguration;
 import org.apache.eventmesh.connector.rocketmq.config.ConfigurationWrapper;
+
 import org.apache.rocketmq.client.exception.MQBrokerException;
 import org.apache.rocketmq.client.exception.MQClientException;
-import org.apache.rocketmq.common.MixAll;
-import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.remoting.exception.RemotingException;
+
+import java.io.File;
+import java.util.Properties;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-public class RocketMQProducerImpl implements MeshMQProducer {
+import io.cloudevents.CloudEvent;
+import io.openmessaging.api.MessagingAccessPoint;
+
+public class RocketMQProducerImpl implements Producer {
 
     public Logger logger = LoggerFactory.getLogger(this.getClass());
 
@@ -67,8 +63,7 @@ public class RocketMQProducerImpl implements MeshMQProducer {
         properties.put("OPERATION_TIMEOUT", 3000);
         properties.put("PRODUCER_ID", producerGroup);
 
-        MessagingAccessPoint messagingAccessPoint = new 
MessagingAccessPointImpl(properties);
-        producer = (ProducerImpl) 
messagingAccessPoint.createProducer(properties);
+        producer = new ProducerImpl(properties);
 
     }
 
@@ -93,20 +88,19 @@ public class RocketMQProducerImpl implements MeshMQProducer 
{
     }
 
     @Override
-    public void send(Message message, SendCallback sendCallback) throws 
Exception {
+    public void publish(CloudEvent message, SendCallback sendCallback) throws 
Exception {
         producer.sendAsync(message, sendCallback);
     }
 
     @Override
-    public void request(Message message, RRCallback rrCallback, long timeout)
+    public void request(CloudEvent message, RRCallback rrCallback, long 
timeout)
             throws InterruptedException, RemotingException, MQClientException, 
MQBrokerException {
         producer.request(message, rrCallback, timeout);
     }
 
     @Override
-    public boolean reply(final Message message, final SendCallback 
sendCallback) throws Exception {
-        message.putSystemProperties(MessageConst.PROPERTY_MESSAGE_TYPE, 
MixAll.REPLY_MESSAGE_FLAG);
-        producer.sendAsync(message, sendCallback);
+    public boolean reply(final CloudEvent message, final SendCallback 
sendCallback) throws Exception {
+        producer.reply(message, sendCallback);
         return true;
     }
 
@@ -121,32 +115,21 @@ public class RocketMQProducerImpl implements 
MeshMQProducer {
     }
 
     @Override
-    public SendResult send(Message message) {
+    public SendResult publish(CloudEvent message) {
         return producer.send(message);
     }
 
     @Override
-    public void sendOneway(Message message) {
+    public void sendOneway(CloudEvent message) {
         producer.sendOneway(message);
     }
 
     @Override
-    public void sendAsync(Message message, SendCallback sendCallback) {
+    public void sendAsync(CloudEvent message, SendCallback sendCallback) {
         producer.sendAsync(message, sendCallback);
     }
 
-    @Override
-    public void setCallbackExecutor(ExecutorService callbackExecutor) {
-        producer.setCallbackExecutor(callbackExecutor);
-    }
 
-    @Override
-    public void updateCredential(Properties credentialProperties) {
-        producer.updateCredential(credentialProperties);
-    }
 
-    @Override
-    public <T> MessageBuilder<T> messageBuilder() {
-        return null;
-    }
+
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/CloudEventUtils.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/CloudEventUtils.java
new file mode 100644
index 0000000..1920ce8
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/eventmesh/connector/rocketmq/utils/CloudEventUtils.java
@@ -0,0 +1,125 @@
+package org.apache.eventmesh.connector.rocketmq.utils;
+
+
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.connector.rocketmq.cloudevent.impl.RocketMQHeaders;
+
+import org.apache.rocketmq.common.message.Message;
+import org.apache.rocketmq.common.message.MessageAccessor;
+import org.apache.rocketmq.common.message.MessageExt;
+
+import java.util.Map;
+import java.util.Set;
+
+public class CloudEventUtils {
+
+    public static SendResult convertSendResult(
+        org.apache.rocketmq.client.producer.SendResult rmqResult) {
+        SendResult sendResult = new SendResult();
+        sendResult.setTopic(rmqResult.getMessageQueue().getTopic());
+        sendResult.setMessageId(rmqResult.getMsgId());
+        return sendResult;
+    }
+
+
+    public static Message msgConvert(MessageExt rmqMsg) {
+        Message message = new Message();
+        if (rmqMsg.getTopic() != null) {
+            message.setTopic(rmqMsg.getTopic());
+        }
+
+        if (rmqMsg.getKeys() != null) {
+            message.setKeys(rmqMsg.getKeys());
+        }
+
+        if (rmqMsg.getTags() != null) {
+            message.setTags(rmqMsg.getTags());
+        }
+
+        if (rmqMsg.getBody() != null) {
+            message.setBody(rmqMsg.getBody());
+        }
+
+        final Set<Map.Entry<String, String>> entries = 
rmqMsg.getProperties().entrySet();
+
+        for (final Map.Entry<String, String> entry : entries) {
+            MessageAccessor.putProperty(message, entry.getKey(), 
entry.getValue());
+        }
+
+        if (rmqMsg.getMsgId() != null) {
+            MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_MESSAGE_ID),
+                rmqMsg.getMsgId());
+        }
+
+        if (rmqMsg.getTopic() != null) {
+            MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_DESTINATION),
+                rmqMsg.getTopic());
+        }
+
+        //
+        
MessageAccessor.putProperty(message,buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_BORN_HOST),
+            String.valueOf(rmqMsg.getBornHost()));
+        MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_BORN_TIMESTAMP),
+            String.valueOf(rmqMsg.getBornTimestamp()));
+        MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_STORE_HOST),
+            String.valueOf(rmqMsg.getStoreHost()));
+        MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey("STORE_TIMESTAMP"),
+            String.valueOf(rmqMsg.getStoreTimestamp()));
+
+        //use in manual ack
+        MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_QUEUE_ID),
+            String.valueOf(rmqMsg.getQueueId()));
+        MessageAccessor.putProperty(message, 
buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_QUEUE_OFFSET),
+            String.valueOf(rmqMsg.getQueueOffset()));
+
+        return message;
+    }
+
+
+
+    private static String buildCloudEventPropertyKey(String propName) {
+        return RocketMQHeaders.CE_PREFIX + propName;
+    }
+
+    public static org.apache.rocketmq.common.message.MessageExt 
msgConvertExt(Message message) {
+
+        org.apache.rocketmq.common.message.MessageExt rmqMessageExt =
+            new org.apache.rocketmq.common.message.MessageExt();
+        try {
+            if (message.getKeys() != null) {
+                rmqMessageExt.setKeys(message.getKeys());
+            }
+            if (message.getTags() != null) {
+                rmqMessageExt.setTags(message.getTags());
+            }
+
+
+            if (message.getBody() != null) {
+                rmqMessageExt.setBody(message.getBody());
+            }
+
+
+            //All destinations in RocketMQ use Topic
+            rmqMessageExt.setTopic(message.getTopic());
+
+            int queueId =
+                (int) 
Integer.valueOf(message.getProperty(buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_QUEUE_ID)));
+            long queueOffset = (long) Long.valueOf(
+                
message.getProperty(buildCloudEventPropertyKey(Constants.PROPERTY_MESSAGE_QUEUE_OFFSET)));
+            //use in manual ack
+            rmqMessageExt.setQueueId(queueId);
+            rmqMessageExt.setQueueOffset(queueOffset);
+            Map<String, String> properties = message.getProperties();
+            for (final Map.Entry<String, String> entry : 
properties.entrySet()) {
+                MessageAccessor.putProperty(rmqMessageExt, entry.getKey(), 
entry.getValue());
+            }
+        } catch (Exception e) {
+            e.printStackTrace();
+        }
+        return rmqMessageExt;
+
+    }
+
+
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/consumer/PushConsumerImplTest.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/consumer/PushConsumerImplTest.java
index 86a5ffd..1a77792 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/consumer/PushConsumerImplTest.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/consumer/PushConsumerImplTest.java
@@ -19,25 +19,19 @@ package org.apache.rocketmq.consumer;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
-import java.lang.reflect.Field;
-import java.util.Collections;
-import java.util.Properties;
-
-import io.openmessaging.api.AsyncConsumeContext;
-import io.openmessaging.api.AsyncMessageListener;
-import io.openmessaging.api.Consumer;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessagingAccessPoint;
-import io.openmessaging.api.OMS;
-import io.openmessaging.api.OMSBuiltinKeys;
-
+import org.apache.eventmesh.api.EventListener;
 import org.apache.eventmesh.api.EventMeshAction;
-import org.apache.eventmesh.api.EventMeshAsyncConsumeContext;
 import org.apache.eventmesh.connector.rocketmq.consumer.PushConsumerImpl;
 import org.apache.eventmesh.connector.rocketmq.domain.NonStandardKeys;
+
 import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
 import 
org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
 import org.apache.rocketmq.common.message.MessageExt;
+
+import java.lang.reflect.Field;
+import java.util.Collections;
+import java.util.Properties;
+
 import org.junit.After;
 import org.junit.Before;
 import org.junit.Test;
@@ -46,9 +40,12 @@ import org.mockito.Mock;
 import org.mockito.Mockito;
 import org.mockito.junit.MockitoJUnitRunner;
 
+import io.cloudevents.CloudEvent;
+import io.openmessaging.api.OMSBuiltinKeys;
+
 @RunWith(MockitoJUnitRunner.class)
 public class PushConsumerImplTest {
-    private Consumer consumer;
+    private PushConsumerImpl consumer;
 
     @Mock
     private DefaultMQPushConsumer rocketmqPushConsumer;
@@ -58,13 +55,13 @@ public class PushConsumerImplTest {
         Properties consumerProp = new Properties();
         consumerProp.setProperty(OMSBuiltinKeys.DRIVER_IMPL, 
"org.apache.eventmesh.connector.rocketmq.MessagingAccessPointImpl");
         consumerProp.setProperty("access_points", "IP1:9876,IP2:9876");
-        final MessagingAccessPoint messagingAccessPoint = 
OMS.builder().build(consumerProp);//.endpoint("oms:rocketmq://IP1:9876,IP2:9876/namespace").build(config);
+        //final MessagingAccessPoint messagingAccessPoint = 
OMS.builder().build(consumerProp);//.endpoint("oms:rocketmq://IP1:9876,IP2:9876/namespace").build(config);
 
         consumerProp.setProperty("message.model", "CLUSTERING");
 
         //Properties consumerProp = new Properties();
         consumerProp.put("CONSUMER_ID", "TestGroup");
-        consumer = messagingAccessPoint.createConsumer(consumerProp);
+        consumer = new PushConsumerImpl(consumerProp);
 
 
         Field field = 
PushConsumerImpl.class.getDeclaredField("rocketmqPushConsumer");
@@ -91,12 +88,14 @@ public class PushConsumerImplTest {
         consumedMsg.setBody(testBody);
         consumedMsg.putUserProperty(NonStandardKeys.MESSAGE_DESTINATION, 
"TOPIC");
         consumedMsg.setTopic("HELLO_QUEUE");
-        consumer.subscribe("HELLO_QUEUE", "*", new AsyncMessageListener() {
+        consumer.subscribe("HELLO_QUEUE", "*", new EventListener() {
+
             @Override
-            public void consume(Message message, AsyncConsumeContext context) {
-                
assertThat(message.getSystemProperties("MESSAGE_ID")).isEqualTo("NewMsgId");
-                assertThat(message.getBody()).isEqualTo(testBody);
-                
((EventMeshAsyncConsumeContext)context).commit(EventMeshAction.CommitMessage);
+            public void consume(CloudEvent cloudEvent,
+                                org.apache.eventmesh.api.AsyncConsumeContext 
context) {
+                
assertThat(cloudEvent.getExtension("MESSAGE_ID")).isEqualTo("NewMsgId");
+                assertThat(cloudEvent.getData()).isEqualTo(testBody);
+                context.commit(EventMeshAction.CommitMessage);
             }
         });
         ((MessageListenerConcurrently) rocketmqPushConsumer
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/producer/ProducerImplTest.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/producer/ProducerImplTest.java
index 83d5571..d4bb646 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/producer/ProducerImplTest.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/test/java/org/apache/rocketmq/producer/ProducerImplTest.java
@@ -21,16 +21,10 @@ import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Fail.failBecauseExceptionWasNotThrown;
 import static org.mockito.ArgumentMatchers.any;
 
-import java.lang.reflect.Field;
-import java.util.Properties;
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+import org.apache.eventmesh.connector.rocketmq.producer.AbstractProducer;
+import org.apache.eventmesh.connector.rocketmq.producer.ProducerImpl;
 
-import io.openmessaging.api.MessagingAccessPoint;
-import io.openmessaging.api.OMS;
-import io.openmessaging.api.OMSBuiltinKeys;
-import io.openmessaging.api.Producer;
-import io.openmessaging.api.exception.OMSRuntimeException;
-
-import org.apache.eventmesh.connector.rocketmq.producer.AbstractOMSProducer;
 import org.apache.rocketmq.client.exception.MQBrokerException;
 import org.apache.rocketmq.client.exception.MQClientException;
 import org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl;
@@ -41,6 +35,11 @@ import org.apache.rocketmq.common.ServiceState;
 import org.apache.rocketmq.common.message.Message;
 import org.apache.rocketmq.common.message.MessageQueue;
 import org.apache.rocketmq.remoting.exception.RemotingException;
+
+import java.lang.reflect.Field;
+import java.net.URI;
+import java.util.Properties;
+
 import org.junit.After;
 import org.junit.Before;
 import org.junit.Test;
@@ -49,9 +48,13 @@ import org.mockito.Mock;
 import org.mockito.Mockito;
 import org.mockito.junit.MockitoJUnitRunner;
 
+import io.cloudevents.CloudEvent;
+import io.cloudevents.core.builder.CloudEventBuilder;
+import io.openmessaging.api.OMSBuiltinKeys;
+
 @RunWith(MockitoJUnitRunner.class)
 public class ProducerImplTest {
-    private Producer producer;
+    private ProducerImpl producer;
 
     @Mock
     private DefaultMQProducer rocketmqProducer;
@@ -61,10 +64,9 @@ public class ProducerImplTest {
         Properties config = new Properties();
         config.setProperty(OMSBuiltinKeys.DRIVER_IMPL, 
"org.apache.eventmesh.connector.rocketmq.MessagingAccessPointImpl");
         config.setProperty("access_points", "IP1:9876,IP2:9876");
-        final MessagingAccessPoint messagingAccessPoint = 
OMS.builder().build(config);//.endpoint("oms:rocketmq://IP1:9876,IP2:9876/namespace").build(config);
-        producer = messagingAccessPoint.createProducer(config);
+        producer = new ProducerImpl(config);
 
-        Field field = 
AbstractOMSProducer.class.getDeclaredField("rocketmqProducer");
+        Field field = 
AbstractProducer.class.getDeclaredField("rocketmqProducer");
         field.setAccessible(true);
         field.set(producer, rocketmqProducer);
 
@@ -97,11 +99,17 @@ public class ProducerImplTest {
         
Mockito.when(rocketmqProducer.getDefaultMQProducerImpl()).thenReturn(defaultMQProducerImpl);
 
 
-        io.openmessaging.api.Message message = new 
io.openmessaging.api.Message("HELLO_TOPIC", "", new byte[]{'a'});
-        io.openmessaging.api.SendResult omsResult =
-                producer.send(message);
+        CloudEvent cloudEvent = CloudEventBuilder.v1()
+            .withId("id1")
+            .withSource(URI.create("https://github.com/cloudevents/*****";))
+            .withType("producer.example")
+            .withSubject("HELLO_TOPIC")
+            .withData("hello world".getBytes())
+            .build();
+        org.apache.eventmesh.api.SendResult result =
+                producer.send(cloudEvent);
 
-        assertThat(omsResult.getMessageId()).isEqualTo("TestMsgID");
+        assertThat(result.getMessageId()).isEqualTo("TestMsgID");
         Mockito.verify(rocketmqProducer).getDefaultMQProducerImpl();
         Mockito.verify(rocketmqProducer).send(any(Message.class));
 
@@ -141,8 +149,15 @@ public class ProducerImplTest {
 
         try {
             io.openmessaging.api.Message message = new 
io.openmessaging.api.Message("HELLO_TOPIC", "", new byte[]{'a'});
-            producer.send(message);
-            failBecauseExceptionWasNotThrown(OMSRuntimeException.class);
+            CloudEvent cloudEvent = CloudEventBuilder.v1()
+                .withId("id1")
+                .withSource(URI.create("https://github.com/cloudevents/*****";))
+                .withType("producer.example")
+                .withSubject("HELLO_TOPIC")
+                .withData(new byte[]{'a'})
+                .build();
+            producer.send(cloudEvent);
+            failBecauseExceptionWasNotThrown(ConnectorRuntimeException.class);
         } catch (Exception e) {
             assertThat(e).hasMessageContaining("Send message to RocketMQ 
broker failed.");
         }

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

Reply via email to