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 5f61c24  Fix standalone connector interface, fix example (#608)
5f61c24 is described below

commit 5f61c248ed4c78eeb4602ebd89b5c4c473aa9269
Author: Wenjun Ruan <[email protected]>
AuthorDate: Wed Nov 24 10:52:59 2021 +0800

    Fix standalone connector interface, fix example (#608)
---
 .../apache/eventmesh/api/consumer/Consumer.java    |   2 -
 .../apache/eventmesh/api/producer/Producer.java    |   6 +-
 .../eventmesh-connector-standalone/build.gradle    |   4 +-
 .../standalone/MessagingAccessPointImpl.java       |  91 -----------
 .../standalone/broker/StandaloneBroker.java        |   9 +-
 .../standalone/broker/model/MessageEntity.java     |   9 +-
 .../standalone/broker/task/SubScribeTask.java      |  27 ++--
 .../standalone/consumer/StandaloneConsumer.java    |  92 ++++-------
 .../consumer/StandaloneConsumerAdaptor.java        | 172 +++++++++++++++++++++
 .../StandaloneMeshMQPushConsumerAdaptor.java       | 135 ----------------
 .../producer/StandaloneMeshMQProducerAdaptor.java  | 125 ---------------
 .../standalone/producer/StandaloneProducer.java    | 131 ++++++++++------
 .../producer/StandaloneProducerAdaptor.java        | 112 ++++++++++++++
 ... => org.apache.eventmesh.api.consumer.Consumer} |   2 +-
 ... => org.apache.eventmesh.api.producer.Producer} |   2 +-
 .../standalone/broker/StandaloneBrokerTest.java    |  11 +-
 .../eventmeshmessage}/AsyncPublishInstance.java    |   9 +-
 .../AsyncSyncRequestInstance.java                  |   2 +-
 .../eventmeshmessage}/SyncRequestInstance.java     |  47 +++---
 .../http/demo/sub/controller/SubController.java    |   7 +-
 .../http/demo/sub/service/SubService.java          |   2 +-
 .../eventmesh/tcp/common/EventMeshTestUtils.java   |  62 ++++----
 .../{ => pub/eventmeshmessage}/AsyncPublish.java   |  24 +--
 .../eventmeshmessage}/AsyncPublishBroadcast.java   |  34 ++--
 .../{ => pub/eventmeshmessage}/SyncRequest.java    |  37 ++---
 .../{ => sub/eventmeshmessage}/AsyncSubscribe.java |  20 ++-
 .../eventmeshmessage}/AsyncSubscribeBroadcast.java |  20 +--
 .../{ => sub/eventmeshmessage}/SyncResponse.java   |  18 ++-
 .../client/http/producer/OpenMessageProducer.java  |   7 +-
 .../eventmesh/client/tcp/common/MessageUtils.java  |  39 ++---
 .../EventMeshMessageTCPPubClient.java              |   3 +-
 .../EventMeshMessageTCPSubClient.java              |   2 +-
 .../client/http/demo/SyncRequestInstance.java      |  45 +++---
 .../client/http/util/HttpLoadBalanceUtilsTest.java |  17 +-
 .../client/tcp/common/EventMeshTestUtils.java      |  63 ++++----
 .../eventmesh/client/tcp/demo/SyncRequest.java     |  25 +--
 36 files changed, 696 insertions(+), 717 deletions(-)

diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
index a6fafe0..87b3ac7 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
@@ -28,8 +28,6 @@ import java.util.Properties;
 
 import io.cloudevents.CloudEvent;
 
-
-
 /**
  * Consumer Interface.
  */
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
index 04662e9..7a9e032 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
@@ -17,7 +17,11 @@
 
 package org.apache.eventmesh.api.producer;
 
-import org.apache.eventmesh.api.*;
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
 import org.apache.eventmesh.spi.EventMeshExtensionType;
 import org.apache.eventmesh.spi.EventMeshSPI;
 
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
index 32b1372..e857140 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
+++ b/eventmesh-connector-plugin/eventmesh-connector-standalone/build.gradle
@@ -16,7 +16,7 @@
  */
 
 dependencies {
-    compileOnly project(":eventmesh-common")
-    compileOnly project(":eventmesh-connector-plugin:eventmesh-connector-api")
+    implementation project(":eventmesh-common")
+    implementation 
project(":eventmesh-connector-plugin:eventmesh-connector-api")
 
 }
\ No newline at end of file
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
deleted file mode 100644
index 6701e81..0000000
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/MessagingAccessPointImpl.java
+++ /dev/null
@@ -1,91 +0,0 @@
-/*
- * 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.standalone;
-
-import io.openmessaging.api.Consumer;
-import io.openmessaging.api.MessagingAccessPoint;
-import io.openmessaging.api.Producer;
-import io.openmessaging.api.PullConsumer;
-import io.openmessaging.api.batch.BatchConsumer;
-import io.openmessaging.api.order.OrderConsumer;
-import io.openmessaging.api.order.OrderProducer;
-import io.openmessaging.api.transaction.LocalTransactionChecker;
-import io.openmessaging.api.transaction.TransactionProducer;
-import org.apache.eventmesh.connector.standalone.consumer.StandaloneConsumer;
-import org.apache.eventmesh.connector.standalone.producer.StandaloneProducer;
-
-import java.util.Properties;
-
-public class MessagingAccessPointImpl implements MessagingAccessPoint {
-
-    private Properties accessPointProperties;
-
-    public MessagingAccessPointImpl(Properties accessPointProperties) {
-        this.accessPointProperties = accessPointProperties;
-    }
-
-    @Override
-    public String version() {
-        return null;
-    }
-
-    @Override
-    public Properties attributes() {
-        return accessPointProperties;
-    }
-
-    @Override
-    public Producer createProducer(Properties properties) {
-        return new StandaloneProducer(properties);
-    }
-
-    @Override
-    public OrderProducer createOrderProducer(Properties properties) {
-        return null;
-    }
-
-    @Override
-    public TransactionProducer createTransactionProducer(Properties 
properties, LocalTransactionChecker checker) {
-        return null;
-    }
-
-    @Override
-    public TransactionProducer createTransactionProducer(Properties 
properties) {
-        return null;
-    }
-
-    @Override
-    public Consumer createConsumer(Properties properties) {
-        return new StandaloneConsumer(properties);
-    }
-
-    @Override
-    public PullConsumer createPullConsumer(Properties properties) {
-        return null;
-    }
-
-    @Override
-    public BatchConsumer createBatchConsumer(Properties properties) {
-        return null;
-    }
-
-    @Override
-    public OrderConsumer createOrderedConsumer(Properties properties) {
-        return null;
-    }
-}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
index 8f35c3a..72a9b58 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBroker.java
@@ -17,6 +17,7 @@
 
 package org.apache.eventmesh.connector.standalone.broker;
 
+import io.cloudevents.CloudEvent;
 import io.openmessaging.api.Message;
 import org.apache.commons.lang3.tuple.Pair;
 import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
@@ -53,7 +54,7 @@ public class StandaloneBroker {
      * @param message   message
      * @throws InterruptedException
      */
-    public MessageEntity putMessage(String topicName, Message message) throws 
InterruptedException {
+    public MessageEntity putMessage(String topicName, CloudEvent message) 
throws InterruptedException {
         Pair<MessageQueue, AtomicLong> pair = createTopicIfAbsent(topicName);
         AtomicLong topicOffset = pair.getRight();
         MessageQueue messageQueue = pair.getLeft();
@@ -70,7 +71,7 @@ public class StandaloneBroker {
      *
      * @param topicName
      */
-    public Message takeMessage(String topicName) throws InterruptedException {
+    public CloudEvent takeMessage(String topicName) throws 
InterruptedException {
         TopicMetadata topicMetadata = new TopicMetadata(topicName);
         return messageContainer.computeIfAbsent(topicMetadata, k -> new 
MessageQueue()).take().getMessage();
     }
@@ -80,7 +81,7 @@ public class StandaloneBroker {
      *
      * @param topicName
      */
-    public Message getMessage(String topicName) {
+    public CloudEvent getMessage(String topicName) {
         TopicMetadata topicMetadata = new TopicMetadata(topicName);
         MessageEntity head = messageContainer.computeIfAbsent(topicMetadata, k 
-> new MessageQueue()).getHead();
         if (head == null) {
@@ -96,7 +97,7 @@ public class StandaloneBroker {
      * @param offset    offset
      * @return
      */
-    public Message getMessage(String topicName, long offset) {
+    public CloudEvent getMessage(String topicName, long offset) {
         TopicMetadata topicMetadata = new TopicMetadata(topicName);
         MessageEntity messageEntity = 
messageContainer.computeIfAbsent(topicMetadata, k -> new 
MessageQueue()).getByOffset(offset);
         if (messageEntity == null) {
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
index d816891..968053a 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/model/MessageEntity.java
@@ -17,6 +17,7 @@
 
 package org.apache.eventmesh.connector.standalone.broker.model;
 
+import io.cloudevents.CloudEvent;
 import io.openmessaging.api.Message;
 
 import java.io.Serializable;
@@ -25,13 +26,13 @@ public class MessageEntity implements Serializable {
 
     private TopicMetadata topicMetadata;
 
-    private Message message;
+    private CloudEvent message;
 
     private long offset;
 
     private long createTimeMills;
 
-    public MessageEntity(TopicMetadata topicMetadata, Message message, long 
offset, long currentTimeMills) {
+    public MessageEntity(TopicMetadata topicMetadata, CloudEvent message, long 
offset, long currentTimeMills) {
         this.topicMetadata = topicMetadata;
         this.message = message;
         this.offset = offset;
@@ -46,11 +47,11 @@ public class MessageEntity implements Serializable {
         this.topicMetadata = topicMetadata;
     }
 
-    public Message getMessage() {
+    public CloudEvent getMessage() {
         return message;
     }
 
-    public void setMessage(Message message) {
+    public void setMessage(CloudEvent message) {
         this.message = message;
     }
 
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
index 430e4aa..bfbf167 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
@@ -27,14 +27,14 @@ import java.util.concurrent.atomic.AtomicInteger;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import io.openmessaging.api.Message;
+import io.cloudevents.CloudEvent;
 
 public class SubScribeTask implements Runnable {
 
-    private String               topicName;
-    private StandaloneBroker     standaloneBroker;
-    private EventListener listener;
-    private volatile boolean isRunning;
+    private          String           topicName;
+    private          StandaloneBroker standaloneBroker;
+    private          EventListener    listener;
+    private volatile boolean          isRunning;
 
     private AtomicInteger offset;
 
@@ -55,13 +55,13 @@ public class SubScribeTask implements Runnable {
             try {
                 logger.debug("execute subscribe task, topic: {}, offset: {}", 
topicName, offset);
                 if (offset == null) {
-                    Message message = standaloneBroker.getMessage(topicName);
+                    CloudEvent message = 
standaloneBroker.getMessage(topicName);
                     if (message != null) {
-                        offset = new AtomicInteger((int) message.getOffset());
+                        offset = new AtomicInteger((int) 
message.getExtension("offset"));
                     }
                 }
                 if (offset != null) {
-                    Message message = standaloneBroker.getMessage(topicName, 
offset.get());
+                    CloudEvent message = 
standaloneBroker.getMessage(topicName, offset.get());
                     if (message != null) {
                         EventMeshAsyncConsumeContext consumeContext = new 
EventMeshAsyncConsumeContext() {
                             @Override
@@ -69,7 +69,8 @@ public class SubScribeTask implements Runnable {
                                 switch (action) {
                                     case CommitMessage:
                                         // update offset
-                                        logger.info("message commit, topic: 
{}, current offset:{}", topicName, offset.get());
+                                        logger.info("message commit, topic: 
{}, current offset:{}", topicName,
+                                            offset.get());
                                         break;
                                     case ReconsumeLater:
                                         // don't update offset
@@ -77,7 +78,8 @@ public class SubScribeTask implements Runnable {
                                     case ManualAck:
                                         // update offset
                                         offset.incrementAndGet();
-                                        logger.info("message ack, topic: {}, 
current offset:{}", topicName, offset.get());
+                                        logger
+                                            .info("message ack, topic: {}, 
current offset:{}", topicName, offset.get());
                                         break;
                                     default:
 
@@ -89,13 +91,14 @@ public class SubScribeTask implements Runnable {
                 }
 
             } catch (Exception ex) {
-                logger.error("consumer error, topic: {}, offset: {}", 
topicName, offset == null ? null : offset.get(), ex);
+                logger.error("consumer error, topic: {}, offset: {}", 
topicName, offset == null ? null : offset.get(),
+                    ex);
             }
             try {
                 Thread.sleep(1000);
             } catch (InterruptedException e) {
                 logger.error("Thread is interrupted, topic: {}, offset: {} 
thread name: {}",
-                        topicName, offset == null ? null : offset.get(), 
Thread.currentThread().getName(), e);
+                    topicName, offset == null ? null : offset.get(), 
Thread.currentThread().getName(), e);
                 Thread.currentThread().interrupt();
             }
         }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
index d206cea..5a18379 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumer.java
@@ -17,23 +17,21 @@
 
 package org.apache.eventmesh.connector.standalone.consumer;
 
-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 org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.consumer.Consumer;
 import org.apache.eventmesh.common.ThreadPoolFactory;
 import org.apache.eventmesh.connector.standalone.broker.StandaloneBroker;
 import org.apache.eventmesh.connector.standalone.broker.model.TopicMetadata;
 import org.apache.eventmesh.connector.standalone.broker.task.SubScribeTask;
 
+import java.util.List;
 import java.util.Properties;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.atomic.AtomicBoolean;
-import java.util.function.Predicate;
+
+import io.cloudevents.CloudEvent;
 
 public class StandaloneConsumer implements Consumer {
 
@@ -50,30 +48,49 @@ public class StandaloneConsumer implements Consumer {
         this.subscribeTaskTable = new ConcurrentHashMap<>(16);
         this.isStarted = new AtomicBoolean(false);
         this.consumeExecutorService = 
ThreadPoolFactory.createThreadPoolExecutor(
-                Runtime.getRuntime().availableProcessors() * 2,
-                Runtime.getRuntime().availableProcessors() * 2,
-                "StandaloneConsumerThread"
+            Runtime.getRuntime().availableProcessors() * 2,
+            Runtime.getRuntime().availableProcessors() * 2,
+            "StandaloneConsumerThread"
         );
     }
 
     @Override
-    public void subscribe(String topic, String subExpression, MessageListener 
listener) {
+    public boolean isStarted() {
+        return isStarted.get();
     }
 
     @Override
-    public void subscribe(String topic, MessageSelector selector, 
MessageListener listener) {
+    public boolean isClosed() {
+        return !isStarted.get();
     }
 
     @Override
-    public <T> void subscribe(String topic, String subExpression, 
GenericMessageListener<T> listener) {
+    public void start() {
+        isStarted.compareAndSet(false, true);
     }
 
     @Override
-    public <T> void subscribe(String topic, MessageSelector selector, 
GenericMessageListener<T> listener) {
+    public void shutdown() {
+        isStarted.compareAndSet(true, false);
+        subscribeTaskTable.forEach(((topic, subScribeTask) -> 
subScribeTask.shutdown()));
+        subscribeTaskTable.clear();
+    }
+
+    @Override
+    public void init(Properties keyValue) throws Exception {
+
     }
 
     @Override
-    public void subscribe(String topic, String subExpression, 
AsyncMessageListener listener) {
+    public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext 
context) {
+        cloudEvents.forEach(cloudEvent -> standaloneBroker.updateOffset(
+            new TopicMetadata(cloudEvent.getSubject()), (Long) 
cloudEvent.getExtension("offset"))
+        );
+
+    }
+
+    @Override
+    public void subscribe(String topic, EventListener listener) throws 
Exception {
         if (listener == null) {
             throw new IllegalArgumentException("listener cannot be null");
         }
@@ -89,18 +106,6 @@ public class StandaloneConsumer implements Consumer {
     }
 
     @Override
-    public void subscribe(String topic, MessageSelector selector, 
AsyncMessageListener listener) {
-    }
-
-    @Override
-    public <T> void subscribe(String topic, String subExpression, 
AsyncGenericMessageListener<T> listener) {
-    }
-
-    @Override
-    public <T> void subscribe(String topic, MessageSelector selector, 
AsyncGenericMessageListener<T> listener) {
-    }
-
-    @Override
     public void unsubscribe(String topic) {
         if (!subscribeTaskTable.containsKey(topic)) {
             return;
@@ -111,35 +116,4 @@ public class StandaloneConsumer implements Consumer {
             subscribeTaskTable.remove(topic);
         }
     }
-
-    @Override
-    public void updateCredential(Properties credentialProperties) {
-
-    }
-
-    @Override
-    public boolean isStarted() {
-        return isStarted.get();
-    }
-
-    @Override
-    public boolean isClosed() {
-        return !isStarted.get();
-    }
-
-    @Override
-    public void start() {
-        isStarted.compareAndSet(false, true);
-    }
-
-    @Override
-    public void shutdown() {
-        isStarted.compareAndSet(true, false);
-        subscribeTaskTable.forEach(((topic, subScribeTask) -> 
subScribeTask.shutdown()));
-        subscribeTaskTable.clear();
-    }
-
-    public void updateOffset(Message message) {
-        standaloneBroker.updateOffset(new TopicMetadata(message.getTopic()), 
message.getOffset());
-    }
 }
\ No newline at end of file
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
new file mode 100644
index 0000000..56518a7
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneConsumerAdaptor.java
@@ -0,0 +1,172 @@
+/*
+ * 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.standalone.consumer;
+
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.consumer.Consumer;
+
+import java.util.List;
+import java.util.Properties;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.cloudevents.CloudEvent;
+
+public class StandaloneConsumerAdaptor implements Consumer {
+
+    private final Logger logger = 
LoggerFactory.getLogger(StandaloneConsumerAdaptor.class);
+
+    private StandaloneConsumer consumer;
+
+    public StandaloneConsumerAdaptor() {
+    }
+
+    @Override
+    public boolean isStarted() {
+        return false;
+    }
+
+    @Override
+    public boolean isClosed() {
+        return false;
+    }
+
+    @Override
+    public void start() {
+
+    }
+
+    @Override
+    public void shutdown() {
+
+    }
+
+    @Override
+    public void init(Properties keyValue) throws Exception {
+
+    }
+
+    @Override
+    public void updateOffset(List<CloudEvent> cloudEvents, AbstractContext 
context) {
+
+    }
+
+    @Override
+    public void subscribe(String topic, EventListener listener) throws 
Exception {
+
+    }
+
+    @Override
+    public void unsubscribe(String topic) {
+
+    }
+
+//    @Override
+//    public void init(Properties keyValue) throws Exception {
+//        String producerGroup = keyValue.getProperty("producerGroup");
+//
+//        MessagingAccessPointImpl messagingAccessPoint = new 
MessagingAccessPointImpl(keyValue);
+//        consumer = (StandaloneConsumer) 
messagingAccessPoint.createConsumer(keyValue);
+//
+//    }
+//
+//    @Override
+//    public void updateOffset(List<Message> msgs, AbstractContext context) {
+//        for(Message message : msgs) {
+//            consumer.updateOffset(message);
+//        }
+//    }
+//
+//    @Override
+//    public void subscribe(String topic, AsyncMessageListener listener) 
throws Exception {
+//        // todo: support subExpression
+//        consumer.subscribe(topic, "*", listener);
+//    }
+//
+//    @Override
+//    public void unsubscribe(String topic) {
+//        consumer.unsubscribe(topic);
+//    }
+//
+//    @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) {
+//
+//    }
+//
+//    @Override
+//    public boolean isStarted() {
+//        return consumer.isStarted();
+//    }
+//
+//    @Override
+//    public boolean isClosed() {
+//        return consumer.isClosed();
+//    }
+//
+//    @Override
+//    public void start() {
+//        consumer.start();
+//    }
+//
+//    @Override
+//    public void shutdown() {
+//        consumer.shutdown();
+//    }
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
deleted file mode 100644
index 3bd4ec9..0000000
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/consumer/StandaloneMeshMQPushConsumerAdaptor.java
+++ /dev/null
@@ -1,135 +0,0 @@
-/*
- * 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.standalone.consumer;
-
-import io.openmessaging.api.AsyncGenericMessageListener;
-import io.openmessaging.api.AsyncMessageListener;
-import io.openmessaging.api.GenericMessageListener;
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageListener;
-import io.openmessaging.api.MessageSelector;
-import org.apache.eventmesh.api.AbstractContext;
-import org.apache.eventmesh.api.consumer.MeshMQPushConsumer;
-import org.apache.eventmesh.connector.standalone.MessagingAccessPointImpl;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.List;
-import java.util.Properties;
-
-public class StandaloneMeshMQPushConsumerAdaptor implements MeshMQPushConsumer 
{
-
-    private final Logger logger = 
LoggerFactory.getLogger(StandaloneMeshMQPushConsumerAdaptor.class);
-
-    private StandaloneConsumer consumer;
-
-    public StandaloneMeshMQPushConsumerAdaptor() {
-    }
-
-    @Override
-    public void init(Properties keyValue) throws Exception {
-        String producerGroup = keyValue.getProperty("producerGroup");
-
-        MessagingAccessPointImpl messagingAccessPoint = new 
MessagingAccessPointImpl(keyValue);
-        consumer = (StandaloneConsumer) 
messagingAccessPoint.createConsumer(keyValue);
-
-    }
-
-    @Override
-    public void updateOffset(List<Message> msgs, AbstractContext context) {
-        for(Message message : msgs) {
-            consumer.updateOffset(message);
-        }
-    }
-
-    @Override
-    public void subscribe(String topic, AsyncMessageListener listener) throws 
Exception {
-        // todo: support subExpression
-        consumer.subscribe(topic, "*", listener);
-    }
-
-    @Override
-    public void unsubscribe(String topic) {
-        consumer.unsubscribe(topic);
-    }
-
-    @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) {
-
-    }
-
-    @Override
-    public boolean isStarted() {
-        return consumer.isStarted();
-    }
-
-    @Override
-    public boolean isClosed() {
-        return consumer.isClosed();
-    }
-
-    @Override
-    public void start() {
-        consumer.start();
-    }
-
-    @Override
-    public void shutdown() {
-        consumer.shutdown();
-    }
-}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
deleted file mode 100644
index fb632c7..0000000
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneMeshMQProducerAdaptor.java
+++ /dev/null
@@ -1,125 +0,0 @@
-/*
- * 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.standalone.producer;
-
-import io.openmessaging.api.Message;
-import io.openmessaging.api.MessageBuilder;
-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.connector.standalone.MessagingAccessPointImpl;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.Properties;
-import java.util.concurrent.ExecutorService;
-
-public class StandaloneMeshMQProducerAdaptor implements MeshMQProducer {
-
-    private final Logger logger = 
LoggerFactory.getLogger(StandaloneMeshMQProducerAdaptor.class);
-
-    private StandaloneProducer standaloneProducer;
-
-    public StandaloneMeshMQProducerAdaptor() {
-    }
-
-    @Override
-    public void init(Properties properties) throws Exception {
-        MessagingAccessPointImpl messagingAccessPoint = new 
MessagingAccessPointImpl(properties);
-        standaloneProducer = (StandaloneProducer) 
messagingAccessPoint.createProducer(properties);
-    }
-
-    @Override
-    public void send(Message message, SendCallback sendCallback) throws 
Exception {
-        standaloneProducer.sendAsync(message, sendCallback);
-    }
-
-    @Override
-    public void request(Message message, RRCallback rrCallback, long timeout) 
throws Exception {
-        throw new UnsupportedOperationException("not support request-reply 
mode when eventstore=standalone");
-    }
-
-    @Override
-    public boolean reply(Message message, SendCallback sendCallback) throws 
Exception {
-        throw new UnsupportedOperationException("not support request-reply 
mode when eventstore=standalone");
-    }
-
-    @Override
-    public void checkTopicExist(String topic) throws Exception {
-        boolean exist = standaloneProducer.checkTopicExist(topic);
-        if (!exist) {
-            throw new RuntimeException(String.format("topic:%s is not exist", 
topic));
-        }
-    }
-
-    @Override
-    public void setExtFields() {
-
-    }
-
-    @Override
-    public SendResult send(Message message) {
-        return standaloneProducer.send(message);
-    }
-
-    @Override
-    public void sendOneway(Message message) {
-        standaloneProducer.sendOneway(message);
-    }
-
-    @Override
-    public void sendAsync(Message message, SendCallback sendCallback) {
-        standaloneProducer.sendAsync(message, sendCallback);
-    }
-
-    @Override
-    public void setCallbackExecutor(ExecutorService callbackExecutor) {
-        standaloneProducer.setCallbackExecutor(callbackExecutor);
-    }
-
-    @Override
-    public void updateCredential(Properties credentialProperties) {
-        standaloneProducer.updateCredential(credentialProperties);
-    }
-
-    @Override
-    public boolean isStarted() {
-        return standaloneProducer.isStarted();
-    }
-
-    @Override
-    public boolean isClosed() {
-        return standaloneProducer.isClosed();
-    }
-
-    @Override
-    public void start() {
-        standaloneProducer.start();
-    }
-
-    @Override
-    public void shutdown() {
-        standaloneProducer.shutdown();
-    }
-
-    @Override
-    public <T> MessageBuilder<T> messageBuilder() {
-        return null;
-    }
-}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
index 1522f51..57c8a2a 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducer.java
@@ -17,23 +17,26 @@
 
 package org.apache.eventmesh.connector.standalone.producer;
 
-import com.google.common.base.Preconditions;
-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 org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.exception.ConnectorRuntimeException;
+import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.api.producer.Producer;
 import org.apache.eventmesh.connector.standalone.broker.StandaloneBroker;
 import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
 
 import java.util.Properties;
-import java.util.concurrent.ExecutorService;
 import java.util.concurrent.atomic.AtomicBoolean;
 
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import com.google.common.base.Preconditions;
+
+import io.cloudevents.CloudEvent;
+
 public class StandaloneProducer implements Producer {
 
     private Logger logger = LoggerFactory.getLogger(StandaloneProducer.class);
@@ -48,80 +51,110 @@ public class StandaloneProducer implements Producer {
     }
 
     @Override
-    public SendResult send(Message message) {
-        Preconditions.checkNotNull(message);
+    public boolean isStarted() {
+        return isStarted.get();
+    }
+
+    @Override
+    public boolean isClosed() {
+        return !isStarted.get();
+    }
+
+    @Override
+    public void start() {
+        isStarted.compareAndSet(false, true);
+    }
+
+    @Override
+    public void shutdown() {
+        isStarted.compareAndSet(true, false);
+    }
+
+    @Override
+    public void init(Properties properties) throws Exception {
+
+    }
+
+    @Override
+    public SendResult publish(CloudEvent cloudEvent) {
+        Preconditions.checkNotNull(cloudEvent);
         try {
-            MessageEntity messageEntity = 
standaloneBroker.putMessage(message.getTopic(), message);
+            MessageEntity messageEntity = 
standaloneBroker.putMessage(cloudEvent.getSubject(), cloudEvent);
             SendResult sendResult = new SendResult();
-            sendResult.setTopic(message.getTopic());
+            sendResult.setTopic(cloudEvent.getSubject());
             sendResult.setMessageId(String.valueOf(messageEntity.getOffset()));
             return sendResult;
         } catch (Exception e) {
-            logger.error("send message error, topic: {}", message.getTopic(), 
e);
-            throw new OMSRuntimeException(String.format("Send message error, 
topic: %s", message.getTopic()));
+            logger.error("send message error, topic: {}", 
cloudEvent.getSubject(), e);
+            throw new ConnectorRuntimeException(
+                String.format("Send message error, topic: %s", 
cloudEvent.getSubject()));
         }
     }
 
     @Override
-    public void sendOneway(Message message) {
-        send(message);
-    }
-
-    @Override
-    public void sendAsync(Message message, SendCallback sendCallback) {
-        Preconditions.checkNotNull(message);
+    public void publish(CloudEvent cloudEvent, SendCallback sendCallback) 
throws Exception {
+        Preconditions.checkNotNull(cloudEvent);
         Preconditions.checkNotNull(sendCallback);
 
         try {
-            SendResult sendResult = send(message);
+            SendResult sendResult = publish(cloudEvent);
             sendCallback.onSuccess(sendResult);
         } catch (Exception ex) {
-            OnExceptionContext exceptionContext = new OnExceptionContext();
-            exceptionContext.setTopic(message.getTopic());
-            exceptionContext.setException(new OMSRuntimeException(ex));
-            exceptionContext.setMessageId(message.getMsgID());
-            sendCallback.onException(exceptionContext);
+            OnExceptionContext onExceptionContext = new OnExceptionContext();
+            onExceptionContext.setMessageId(cloudEvent.getId());
+            onExceptionContext.setTopic(cloudEvent.getSubject());
+            onExceptionContext.setException(new ConnectorRuntimeException(ex));
+            sendCallback.onException(onExceptionContext);
         }
-
     }
 
     @Override
-    public void setCallbackExecutor(ExecutorService callbackExecutor) {
-
+    public void sendOneway(CloudEvent cloudEvent) {
+        publish(cloudEvent);
     }
 
     @Override
-    public void updateCredential(Properties credentialProperties) {
-
+    public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
+        Preconditions.checkNotNull(cloudEvent);
+        Preconditions.checkNotNull(sendCallback);
+        // todo: current is not async
+        try {
+            SendResult sendResult = publish(cloudEvent);
+            sendCallback.onSuccess(sendResult);
+        } catch (Exception ex) {
+            OnExceptionContext onExceptionContext = new OnExceptionContext();
+            onExceptionContext.setMessageId(cloudEvent.getId());
+            onExceptionContext.setTopic(cloudEvent.getSubject());
+            onExceptionContext.setException(new ConnectorRuntimeException(ex));
+            sendCallback.onException(onExceptionContext);
+        }
     }
 
     @Override
-    public boolean isStarted() {
-        return isStarted.get();
+    public void request(CloudEvent cloudEvent, RRCallback rrCallback, long 
timeout) throws Exception {
+        throw new ConnectorRuntimeException("Request is not supported");
     }
 
     @Override
-    public boolean isClosed() {
-        return !isStarted.get();
+    public void request(CloudEvent cloudEvent, RequestReplyCallback 
rrCallback, long timeout) throws Exception {
+        throw new ConnectorRuntimeException("Request is not supported");
     }
 
     @Override
-    public void start() {
-        isStarted.compareAndSet(false, true);
+    public boolean reply(CloudEvent cloudEvent, SendCallback sendCallback) 
throws Exception {
+        throw new ConnectorRuntimeException("Reply is not supported");
     }
 
     @Override
-    public void shutdown() {
-        isStarted.compareAndSet(true, false);
-
+    public void checkTopicExist(String topic) throws Exception {
+        boolean exist = standaloneBroker.checkTopicExist(topic);
+        if (!exist) {
+            throw new ConnectorRuntimeException(String.format("topic:%s is not 
exist", topic));
+        }
     }
 
     @Override
-    public <T> MessageBuilder<T> messageBuilder() {
-        return null;
-    }
+    public void setExtFields() {
 
-    public boolean checkTopicExist(String topicName) {
-        return standaloneBroker.checkTopicExist(topicName);
     }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
new file mode 100644
index 0000000..d60968c
--- /dev/null
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/producer/StandaloneProducerAdaptor.java
@@ -0,0 +1,112 @@
+/*
+ * 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.standalone.producer;
+
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.RequestReplyCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.api.producer.Producer;
+
+import java.util.Properties;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import io.cloudevents.CloudEvent;
+
+public class StandaloneProducerAdaptor implements Producer {
+
+    private final Logger logger = 
LoggerFactory.getLogger(StandaloneProducerAdaptor.class);
+
+    private StandaloneProducer standaloneProducer;
+
+    public StandaloneProducerAdaptor() {
+    }
+
+    @Override
+    public boolean isStarted() {
+        return standaloneProducer.isStarted();
+    }
+
+    @Override
+    public boolean isClosed() {
+        return standaloneProducer.isClosed();
+    }
+
+    @Override
+    public void start() {
+        standaloneProducer.start();
+    }
+
+    @Override
+    public void shutdown() {
+        standaloneProducer.shutdown();
+    }
+
+    @Override
+    public void init(Properties properties) throws Exception {
+        standaloneProducer.init(properties);
+    }
+
+    @Override
+    public SendResult publish(CloudEvent cloudEvent) {
+        return standaloneProducer.publish(cloudEvent);
+    }
+
+    @Override
+    public void publish(CloudEvent cloudEvent, SendCallback sendCallback) 
throws Exception {
+        standaloneProducer.publish(cloudEvent, sendCallback);
+    }
+
+    @Override
+    public void sendOneway(CloudEvent cloudEvent) {
+        standaloneProducer.sendOneway(cloudEvent);
+    }
+
+    @Override
+    public void sendAsync(CloudEvent cloudEvent, SendCallback sendCallback) {
+        standaloneProducer.sendAsync(cloudEvent, sendCallback);
+    }
+
+    @Override
+    public void request(CloudEvent cloudEvent, RRCallback rrCallback, long 
timeout) throws Exception {
+        standaloneProducer.request(cloudEvent, rrCallback, timeout);
+    }
+
+    @Override
+    public void request(CloudEvent cloudEvent, RequestReplyCallback 
rrCallback, long timeout) throws Exception {
+        standaloneProducer.request(cloudEvent, rrCallback, timeout);
+    }
+
+    @Override
+    public boolean reply(CloudEvent cloudEvent, SendCallback sendCallback) 
throws Exception {
+        return standaloneProducer.reply(cloudEvent, sendCallback);
+    }
+
+    @Override
+    public void checkTopicExist(String topic) throws Exception {
+        standaloneProducer.checkTopicExist(topic);
+    }
+
+    @Override
+    public void setExtFields() {
+        standaloneProducer.setExtFields();
+    }
+
+}
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
similarity index 96%
rename from 
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
rename to 
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
index 190a92a..5ceb0ed 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.MeshMQPushConsumer
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.consumer.Consumer
@@ -13,4 +13,4 @@
 # See the License for the specific language governing permissions and
 # limitations under the License.
 
-standalone=org.apache.eventmesh.connector.standalone.consumer.StandaloneMeshMQPushConsumerAdaptor
\ No newline at end of file
+standalone=org.apache.eventmesh.connector.standalone.consumer.StandaloneConsumerAdaptor
\ No newline at end of file
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
similarity index 96%
rename from 
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
rename to 
eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
index 4bbac31..1b30d98 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.MeshMQProducer
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/resources/META-INF/eventmesh/org.apache.eventmesh.api.producer.Producer
@@ -13,4 +13,4 @@
 # See the License for the specific language governing permissions and
 # limitations under the License.
 
-standalone=org.apache.eventmesh.connector.standalone.producer.StandaloneMeshMQProducerAdaptor
\ No newline at end of file
+standalone=org.apache.eventmesh.connector.standalone.producer.StandaloneProducerAdaptor
\ No newline at end of file
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
index a63be75..92e4eb7 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/test/java/org/apache/eventmesh/connector/standalone/broker/StandaloneBrokerTest.java
@@ -17,11 +17,14 @@
 
 package org.apache.eventmesh.connector.standalone.broker;
 
-import io.openmessaging.api.Message;
 import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
+
 import org.junit.Assert;
 import org.junit.Test;
 
+import io.cloudevents.CloudEvent;
+import io.cloudevents.core.builder.CloudEventBuilder;
+
 public class StandaloneBrokerTest {
 
     @Test
@@ -32,15 +35,15 @@ public class StandaloneBrokerTest {
     @Test
     public void putMessage() throws InterruptedException {
         StandaloneBroker instance = StandaloneBroker.getInstance();
-        MessageEntity messageEntity = instance.putMessage("test-topic", new 
Message());
+        MessageEntity messageEntity = instance.putMessage("test-topic", 
CloudEventBuilder.v1().build());
         Assert.assertNotNull(messageEntity);
     }
 
     @Test
     public void takeMessage() throws InterruptedException {
         StandaloneBroker instance = StandaloneBroker.getInstance();
-        instance.putMessage("test-topic", new Message());
-        Message message = instance.takeMessage("test-topic");
+        instance.putMessage("test-topic", CloudEventBuilder.v1().build());
+        CloudEvent message = instance.takeMessage("test-topic");
         Assert.assertNotNull(message);
     }
 
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
similarity index 95%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
index 1af175d..8c90d8e 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
@@ -15,7 +15,9 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.http.demo;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
+
+import lombok.extern.slf4j.Slf4j;
 
 import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
 import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
@@ -33,10 +35,9 @@ import java.util.Properties;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+@Slf4j
 public class AsyncPublishInstance {
 
-    public static Logger logger = 
LoggerFactory.getLogger(AsyncPublishInstance.class);
-
     // This messageSize is also used in SubService.java (Subscriber)
     public static int messageSize = 5;
 
@@ -45,7 +46,7 @@ public class AsyncPublishInstance {
         final String eventMeshIp = properties.getProperty("eventmesh.ip");
         final String eventMeshHttpPort = 
properties.getProperty("eventmesh.http.port");
 
-        String eventMeshIPPort;
+        final String eventMeshIPPort;
         if (StringUtils.isBlank(eventMeshIp) || 
StringUtils.isBlank(eventMeshHttpPort)) {
             // if has multi value, can config as: 
127.0.0.1:10105;127.0.0.2:10105
             eventMeshIPPort = "127.0.0.1:10105";
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
similarity index 98%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
index a007ed1..4284726 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/AsyncSyncRequestInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.http.demo;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
 
 import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
 import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
similarity index 70%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
index 58d287d..4d7e1d2 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/SyncRequestInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
@@ -15,16 +15,17 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.http.demo;
-
-import org.apache.commons.lang3.StringUtils;
+package org.apache.eventmesh.http.demo.pub.eventmeshmessage;
 
 import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
 import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
-import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
 import org.apache.eventmesh.common.utils.ThreadUtils;
+
+import org.apache.commons.lang3.StringUtils;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -50,38 +51,36 @@ public class SyncRequestInstance {
                 eventMeshIPPort = "127.0.0.1:10105";
             }
 
-            EventMeshHttpClientConfig eventMeshClientConfig = new 
EventMeshHttpClientConfig();
-            eventMeshClientConfig.setLiteEventMeshAddr(eventMeshIPPort)
-                    .setProducerGroup("EventMeshTest-producerGroup")
-                    .setEnv("env")
-                    .setIdc("idc")
-                    .setIp(IPUtils.getLocalAddress())
-                    .setSys("1234")
-                    .setPid(String.valueOf(ThreadUtils.getPID()));
+            EventMeshHttpClientConfig eventMeshClientConfig = 
EventMeshHttpClientConfig.builder()
+                .liteEventMeshAddr(eventMeshIPPort)
+                .producerGroup("EventMeshTest-producerGroup")
+                .env("env")
+                .idc("idc")
+                .ip(IPUtils.getLocalAddress())
+                .sys("1234")
+                .pid(String.valueOf(ThreadUtils.getPID())).build();
 
             eventMeshHttpProducer = new 
EventMeshHttpProducer(eventMeshClientConfig);
-            eventMeshHttpProducer.start();
 
             long startTime = System.currentTimeMillis();
-            EventMeshMessage eventMeshMessage = new EventMeshMessage();
-            eventMeshMessage.setBizSeqNo(RandomStringUtils.generateNum(30))
-                    .setContent("contentStr with special protocal")
-                    .setTopic(topic)
-                    .setUniqueId(RandomStringUtils.generateNum(30));
+            EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+                .bizSeqNo(RandomStringUtils.generateNum(30))
+                .content("contentStr with special protocal")
+                .topic(topic)
+                .uniqueId(RandomStringUtils.generateNum(30)).build();
 
             EventMeshMessage rsp = 
eventMeshHttpProducer.request(eventMeshMessage, 10000);
             if (logger.isDebugEnabled()) {
-                logger.debug("sendmsg : {}, return : {}, cost:{}ms", 
eventMeshMessage.getContent(), rsp.getContent(), System.currentTimeMillis() - 
startTime);
+                logger.debug("sendmsg : {}, return : {}, cost:{}ms", 
eventMeshMessage.getContent(), rsp.getContent(),
+                    System.currentTimeMillis() - startTime);
             }
         } catch (Exception e) {
             logger.warn("send msg failed", e);
         }
 
-        try {
-            Thread.sleep(30000);
-            if (eventMeshHttpProducer != null) {
-                eventMeshHttpProducer.shutdown();
-            }
+        Thread.sleep(30000);
+        try (final EventMeshHttpProducer close = eventMeshHttpProducer) {
+            // close producer
         } catch (Exception e1) {
             logger.warn("producer shutdown exception", e1);
         }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
index aa400ce..2f0089c 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/controller/SubController.java
@@ -17,6 +17,8 @@
 
 package org.apache.eventmesh.http.demo.sub.controller;
 
+import lombok.extern.slf4j.Slf4j;
+
 import org.apache.eventmesh.common.utils.JsonUtils;
 import org.apache.eventmesh.http.demo.sub.service.SubService;
 
@@ -31,18 +33,17 @@ import 
org.springframework.web.bind.annotation.RequestMapping;
 import org.springframework.web.bind.annotation.RequestMethod;
 import org.springframework.web.bind.annotation.RestController;
 
+@Slf4j
 @RestController
 @RequestMapping("/sub")
 public class SubController {
 
-    public static Logger logger = LoggerFactory.getLogger(SubController.class);
-
     @Autowired
     private SubService subService;
 
     @RequestMapping(value = "/test", method = RequestMethod.POST)
     public String subTest(@RequestBody String message) {
-        logger.info("=======receive message======= {}", 
JsonUtils.serialize(message));
+        log.info("=======receive message======= {}", 
JsonUtils.serialize(message));
         subService.consumeMessage(message);
 
         Map<String, Object> map = new HashMap<>();
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
index 6470c77..0c0f2c1 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/sub/service/SubService.java
@@ -26,7 +26,7 @@ import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
 import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.utils.ThreadUtils;
-import org.apache.eventmesh.http.demo.AsyncPublishInstance;
+import 
org.apache.eventmesh.http.demo.pub.eventmeshmessage.AsyncPublishInstance;
 import org.apache.eventmesh.util.Utils;
 
 import java.util.ArrayList;
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
index 854eb5b..bd7f7cb 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/common/EventMeshTestUtils.java
@@ -22,49 +22,49 @@ import static 
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PR
 import static 
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientBroadCast;
 import static 
org.apache.eventmesh.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientUniCast;
 
-import java.util.concurrent.ThreadLocalRandom;
-
 import org.apache.eventmesh.common.protocol.tcp.Command;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Header;
 import org.apache.eventmesh.common.protocol.tcp.Package;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 
+import java.util.concurrent.ThreadLocalRandom;
+
 public class EventMeshTestUtils {
     private static final int seqLength = 10;
 
     public static UserAgent generateClient1() {
         return UserAgent.builder()
-                .env("test")
-                .host("127.0.0.1")
-                .password(generateRandomString(8))
-                .username("PU4283")
-                .producerGroup("EventmeshTest-ProducerGroup")
-                .consumerGroup("EventmeshTest-ConsumerGroup")
-                .path("/data/app/umg_proxy")
-                .port(8362)
-                .subsystem("5023")
-                .pid(32893)
-                .version("2.0.11")
-                .idc("FT")
-                .build();
+            .env("test")
+            .host("127.0.0.1")
+            .password(generateRandomString(8))
+            .username("PU4283")
+            .producerGroup("EventmeshTest-ProducerGroup")
+            .consumerGroup("EventmeshTest-ConsumerGroup")
+            .path("/data/app/umg_proxy")
+            .port(8362)
+            .subsystem("5023")
+            .pid(32893)
+            .version("2.0.11")
+            .idc("FT")
+            .build();
     }
 
     public static UserAgent generateClient2() {
         return UserAgent.builder()
-                .env("test")
-                .host("127.0.0.1")
-                .password(generateRandomString(8))
-                .username("PU4283")
-                .producerGroup("EventmeshTest-ProducerGroup")
-                .consumerGroup("EventmeshTest-ConsumerGroup")
-                .path("/data/app/umg_proxy")
-                .port(9362)
-                .subsystem("5017")
-                .pid(42893)
-                .version("2.0.11")
-                .idc("FT")
-                .build();
+            .env("test")
+            .host("127.0.0.1")
+            .password(generateRandomString(8))
+            .username("PU4283")
+            .producerGroup("EventmeshTest-ProducerGroup")
+            .consumerGroup("EventmeshTest-ConsumerGroup")
+            .path("/data/app/umg_proxy")
+            .port(9362)
+            .subsystem("5017")
+            .pid(42893)
+            .version("2.0.11")
+            .idc("FT")
+            .build();
     }
 
     public static Package syncRR() {
@@ -102,7 +102,7 @@ public class EventMeshTestUtils {
         return msg;
     }
 
-    private static EventMeshMessage generateSyncRRMqMsg() {
+    public static EventMeshMessage generateSyncRRMqMsg() {
         EventMeshMessage mqMsg = new EventMeshMessage();
         mqMsg.setTopic(TOPIC_PRX_SyncSubscribeTest);
         mqMsg.getProperties().put("msgType", "persistent");
@@ -123,7 +123,7 @@ public class EventMeshTestUtils {
         return mqMsg;
     }
 
-    private static EventMeshMessage generateAsyncEventMqMsg() {
+    public static EventMeshMessage generateAsyncEventMqMsg() {
         EventMeshMessage mqMsg = new EventMeshMessage();
         mqMsg.setTopic(TOPIC_PRX_WQ2ClientUniCast);
         mqMsg.getProperties().put("REPLY_TO", 
"10.36.0.109@ProducerGroup-producerPool-9-access#V1_4_0#CI");
@@ -133,7 +133,7 @@ public class EventMeshTestUtils {
         return mqMsg;
     }
 
-    private static EventMeshMessage generateBroadcastMqMsg() {
+    public static EventMeshMessage generateBroadcastMqMsg() {
         EventMeshMessage mqMsg = new EventMeshMessage();
         mqMsg.setTopic(TOPIC_PRX_WQ2ClientBroadCast);
         mqMsg.getProperties().put("REPLY_TO", 
"10.36.0.109@ProducerGroup-producerPool-9-access#V1_4_0#CI");
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
similarity index 70%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
index 90bd4c7..52ccc15 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublish.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
@@ -15,14 +15,14 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
 
 import java.util.Properties;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 import org.apache.eventmesh.util.Utils;
@@ -33,7 +33,7 @@ public class AsyncPublish {
 
     public static Logger logger = LoggerFactory.getLogger(AsyncPublish.class);
 
-    private static EventMeshTCPClient client;
+    private static EventMeshMessageTCPPubClient client;
 
     public static AsyncPublish handler = new AsyncPublish();
 
@@ -43,14 +43,20 @@ public class AsyncPublish {
         final int eventMeshTcpPort = 
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
         try {
             UserAgent userAgent = EventMeshTestUtils.generateClient1();
-            client = new DefaultEventMeshTCPClient(eventMeshIp, 
eventMeshTcpPort, userAgent);
+            EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+                .host(eventMeshIp)
+                .port(eventMeshTcpPort)
+                .userAgent(userAgent)
+                .build();
+            client = new 
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig);
             client.init();
             client.heartbeat();
 
             for (int i = 0; i < 5; i++) {
-                Package asyncMsg = EventMeshTestUtils.asyncMessage();
-                logger.info("begin send async msg[{}]==================={}", 
i, asyncMsg);
-                client.publish(asyncMsg, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+                EventMeshMessage eventMeshMessage = 
EventMeshTestUtils.generateAsyncEventMqMsg();
+
+                logger.info("begin send async msg[{}]==================={}", 
i, eventMeshMessage);
+                client.publish(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
 
                 Thread.sleep(1000);
             }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
similarity index 67%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
index f1dadef..12071de 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncPublishBroadcast.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
@@ -15,17 +15,18 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
 
-import java.util.Properties;
-
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 import org.apache.eventmesh.util.Utils;
+
+import java.util.Properties;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -33,25 +34,26 @@ public class AsyncPublishBroadcast {
 
     public static Logger logger = 
LoggerFactory.getLogger(AsyncPublishBroadcast.class);
 
-    private static EventMeshTCPClient client;
-
     public static void main(String[] agrs) throws Exception {
         Properties properties = 
Utils.readPropertiesFile("application.properties");
         final String eventMeshIp = properties.getProperty("eventmesh.ip");
         final int eventMeshTcpPort = 
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
-        try {
-            UserAgent userAgent = EventMeshTestUtils.generateClient1();
-            client = new DefaultEventMeshTCPClient(eventMeshIp, 
eventMeshTcpPort, userAgent);
+        UserAgent userAgent = EventMeshTestUtils.generateClient1();
+        EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+            .host(eventMeshIp)
+            .port(eventMeshTcpPort)
+            .userAgent(userAgent)
+            .build();
+        try (final EventMeshMessageTCPPubClient client = new 
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig)) {
             client.init();
             client.heartbeat();
 
-            Package broadcastMsg = EventMeshTestUtils.broadcastMessage();
-            logger.info("begin send broadcast msg============={}", 
broadcastMsg);
-            client.broadcast(broadcastMsg, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+            EventMeshMessage eventMeshMessage = 
EventMeshTestUtils.generateBroadcastMqMsg();
+            logger.info("begin send broadcast msg============={}", 
eventMeshMessage);
+            client.broadcast(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
 
             Thread.sleep(2000);
-            // release resource and close client
-            // client.close();
+
         } catch (Exception e) {
             logger.warn("AsyncPublishBroadcast failed", e);
         }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
similarity index 51%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
index 65339de..844b9cf 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncRequest.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/SyncRequest.java
@@ -15,40 +15,41 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.pub.eventmeshmessage;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Package;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import lombok.extern.slf4j.Slf4j;
 
+@Slf4j
 public class SyncRequest {
 
-    public static Logger logger = LoggerFactory.getLogger(SyncRequest.class);
-
-    private static EventMeshTCPClient client;
+    private static EventMeshMessageTCPPubClient client;
 
     public static void main(String[] agrs) throws Exception {
-        try {
-            UserAgent userAgent = EventMeshTestUtils.generateClient1();
-            client = new DefaultEventMeshTCPClient("127.0.0.1", 10000, 
userAgent);
+        UserAgent userAgent = EventMeshTestUtils.generateClient1();
+        EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+            .host("127.0.0.1")
+            .port(10000)
+            .userAgent(userAgent)
+            .build();
+        try (EventMeshMessageTCPPubClient client = new 
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig)) {
             client.init();
             client.heartbeat();
 
-            Package rrMsg = EventMeshTestUtils.syncRR();
-            logger.info("begin send rr msg=================={}", rrMsg);
-            Package response = client.rr(rrMsg, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
-            logger.info("receive rr reply==================={}", response);
+            EventMeshMessage eventMeshMessage = 
EventMeshTestUtils.generateSyncRRMqMsg();
+            log.info("begin send rr msg=================={}", 
eventMeshMessage);
+            Package response = client.rr(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+            log.info("receive rr reply==================={}", response);
 
-            // release resource and close client
-            // client.close();
         } catch (Exception e) {
-            logger.warn("SyncRequest failed", e);
+            log.warn("SyncRequest failed", e);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
similarity index 77%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
index 18d607a..ec56aa5 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribe.java
@@ -15,11 +15,11 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,7 +36,7 @@ import lombok.extern.slf4j.Slf4j;
 @Slf4j
 public class AsyncSubscribe implements ReceiveMsgHook<EventMeshMessage> {
 
-    private static EventMeshTCPClient client;
+    private static EventMeshMessageTCPSubClient client;
 
     public static AsyncSubscribe handler = new AsyncSubscribe();
 
@@ -44,14 +44,18 @@ public class AsyncSubscribe implements 
ReceiveMsgHook<EventMeshMessage> {
         Properties properties = 
Utils.readPropertiesFile("application.properties");
         final String eventMeshIp = properties.getProperty("eventmesh.ip");
         final int eventMeshTcpPort = 
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
-        try {
-            UserAgent userAgent = EventMeshTestUtils.generateClient2();
-            client = new DefaultEventMeshTCPClient(eventMeshIp, 
eventMeshTcpPort, userAgent);
+        UserAgent userAgent = EventMeshTestUtils.generateClient2();
+        EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+            .host(eventMeshIp)
+            .port(eventMeshTcpPort)
+            .userAgent(userAgent)
+            .build();
+        try (EventMeshMessageTCPSubClient client = new 
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
             client.init();
             client.heartbeat();
 
             client.subscribe("TEST-TOPIC-TCP-ASYNC", 
SubscriptionMode.CLUSTERING, SubscriptionType.ASYNC);
-            client.registerSubBusiHandler(handler);
+            client.registerBusiHandler(handler);
 
             client.listen();
 
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
similarity index 79%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
index 1d6c6d1..8d642e1 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/AsyncSubscribeBroadcast.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/AsyncSubscribeBroadcast.java
@@ -15,11 +15,11 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,22 +36,24 @@ import lombok.extern.slf4j.Slf4j;
 @Slf4j
 public class AsyncSubscribeBroadcast implements 
ReceiveMsgHook<EventMeshMessage> {
 
-    private static EventMeshTCPClient client;
-
     public static AsyncSubscribeBroadcast handler = new 
AsyncSubscribeBroadcast();
 
     public static void main(String[] agrs) throws Exception {
         Properties properties = 
Utils.readPropertiesFile("application.properties");
         final String eventMeshIp = properties.getProperty("eventmesh.ip");
         final int eventMeshTcpPort = 
Integer.parseInt(properties.getProperty("eventmesh.tcp.port"));
-        try {
-            UserAgent userAgent = EventMeshTestUtils.generateClient2();
-            client = new DefaultEventMeshTCPClient(eventMeshIp, 
eventMeshTcpPort, userAgent);
+        UserAgent userAgent = EventMeshTestUtils.generateClient2();
+        EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+            .host(eventMeshIp)
+            .port(eventMeshTcpPort)
+            .userAgent(userAgent)
+            .build();
+        try (EventMeshMessageTCPSubClient client = new 
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
             client.init();
             client.heartbeat();
 
             client.subscribe("TEST-TOPIC-TCP-BROADCAST", 
SubscriptionMode.BROADCASTING, SubscriptionType.ASYNC);
-            client.registerSubBusiHandler(handler);
+            client.registerBusiHandler(handler);
 
             client.listen();
 
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
similarity index 76%
rename from 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
rename to 
eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
index 026fa11..b796c4b 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/SyncResponse.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/sub/eventmeshmessage/SyncResponse.java
@@ -15,11 +15,11 @@
  * limitations under the License.
  */
 
-package org.apache.eventmesh.tcp.demo;
+package org.apache.eventmesh.tcp.demo.sub.eventmeshmessage;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPSubClient;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
@@ -36,15 +36,19 @@ public class SyncResponse implements 
ReceiveMsgHook<EventMeshMessage> {
     public static SyncResponse handler = new SyncResponse();
 
     public static void main(String[] agrs) throws Exception {
-        try {
-            UserAgent userAgent = EventMeshTestUtils.generateClient2();
-            EventMeshTCPClient client = new 
DefaultEventMeshTCPClient("127.0.0.1", 10000, userAgent);
+        UserAgent userAgent = EventMeshTestUtils.generateClient2();
+        EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+            .host("127.0.0.1")
+            .port(10000)
+            .userAgent(userAgent)
+            .build();
+        try (EventMeshMessageTCPSubClient client = new 
EventMeshMessageTCPSubClient(eventMeshTcpClientConfig)) {
             client.init();
             client.heartbeat();
 
             client.subscribe("TEST-TOPIC-TCP-SYNC", 
SubscriptionMode.CLUSTERING, SubscriptionType.SYNC);
             // Synchronize RR messages
-            client.registerSubBusiHandler(handler);
+            client.registerBusiHandler(handler);
 
             client.listen();
 
diff --git 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
index b3cadcc..4dcffff 100644
--- 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
+++ 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/http/producer/OpenMessageProducer.java
@@ -26,6 +26,10 @@ import lombok.extern.slf4j.Slf4j;
 @Slf4j
 class OpenMessageProducer extends AbstractHttpClient implements 
EventMeshProtocolProducer<Message> {
 
+    private static final String PROTOCOL_TYPE = "openmessage";
+
+    private static final String PROTOCOL_DESC = "http";
+
     public OpenMessageProducer(EventMeshHttpClientConfig 
eventMeshHttpClientConfig)
         throws EventMeshException {
         super(eventMeshHttpClientConfig);
@@ -94,8 +98,9 @@ class OpenMessageProducer extends AbstractHttpClient 
implements EventMeshProtoco
         requestParam
             .addHeader(ProtocolKey.ClientInstanceKey.USERNAME, 
eventMeshHttpClientConfig.getUserName())
             .addHeader(ProtocolKey.ClientInstanceKey.PASSWD, 
eventMeshHttpClientConfig.getPassword())
-            .addHeader(ProtocolKey.VERSION, ProtocolVersion.V1.getVersion())
             .addHeader(ProtocolKey.LANGUAGE, Constants.LANGUAGE_JAVA)
+            .addHeader(ProtocolKey.PROTOCOL_TYPE, PROTOCOL_TYPE)
+            .addHeader(ProtocolKey.PROTOCOL_DESC, PROTOCOL_DESC)
             // todo: add producerGroup to header, set protocol type, protocol 
version
             .addBody(SendMessageRequestBody.PRODUCERGROUP, 
eventMeshHttpClientConfig.getProducerGroup())
             .addBody(SendMessageRequestBody.CONTENT, 
JsonUtils.serialize(openMessage));
diff --git 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
index e799fb7..0b96cab 100644
--- 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
+++ 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/MessageUtils.java
@@ -17,19 +17,23 @@
 
 package org.apache.eventmesh.client.tcp.common;
 
+import org.apache.eventmesh.common.Constants;
+import org.apache.eventmesh.common.protocol.SubscriptionItem;
+import org.apache.eventmesh.common.protocol.SubscriptionMode;
+import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
+import org.apache.eventmesh.common.protocol.tcp.Header;
+import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.Subscription;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+
 import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.ThreadLocalRandom;
 
 import io.cloudevents.CloudEvent;
 import io.cloudevents.SpecVersion;
-import org.apache.eventmesh.common.Constants;
-import org.apache.eventmesh.common.protocol.SubscriptionType;
-import org.apache.eventmesh.common.protocol.tcp.Subscription;
-import org.apache.eventmesh.common.protocol.SubscriptionItem;
-import org.apache.eventmesh.common.protocol.SubscriptionMode;
-import org.apache.eventmesh.common.protocol.tcp.*;
-import org.apache.eventmesh.common.protocol.tcp.Package;
 
 public class MessageUtils {
     private static final int seqLength = 10;
@@ -59,7 +63,8 @@ public class MessageUtils {
         return msg;
     }
 
-    public static Package subscribe(String topic, SubscriptionMode 
subscriptionMode, SubscriptionType subscriptionType) {
+    public static Package subscribe(String topic, SubscriptionMode 
subscriptionMode,
+                                    SubscriptionType subscriptionType) {
         Package msg = new Package();
         msg.setHeader(new Header(Command.SUBSCRIBE_REQUEST, 0, null, 
generateRandomString(seqLength)));
         msg.setBody(generateSubscription(topic, subscriptionMode, 
subscriptionType));
@@ -81,18 +86,13 @@ public class MessageUtils {
 
     public static Package buildPackage(Object message, Command command) {
         Package msg = new Package();
-        msg.setHeader(new Header(command, 0,
-                null, generateRandomString(seqLength)));
+        msg.setHeader(new Header(command, 0, null, 
generateRandomString(seqLength)));
         if (message instanceof CloudEvent) {
-            msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
-                    EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
-            msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
-                    ((CloudEvent) message).getSpecVersion().toString());
+            msg.getHeader().putProperty(Constants.PROTOCOL_TYPE, 
EventMeshCommon.CLOUD_EVENTS_PROTOCOL_NAME);
+            msg.getHeader().putProperty(Constants.PROTOCOL_VERSION, 
((CloudEvent) message).getSpecVersion().toString());
         } else if (message instanceof EventMeshMessage) {
-            msg.getHeader().putProperty(Constants.PROTOCOL_TYPE,
-                    EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME);
-            msg.getHeader().putProperty(Constants.PROTOCOL_VERSION,
-                    SpecVersion.V1.toString());
+            msg.getHeader().putProperty(Constants.PROTOCOL_TYPE, 
EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME);
+            msg.getHeader().putProperty(Constants.PROTOCOL_VERSION, 
SpecVersion.V1.toString());
         } else {
             // unsupported protocol for server
             return msg;
@@ -159,7 +159,8 @@ public class MessageUtils {
             .build();
     }
 
-    private static Subscription generateSubscription(String topic, 
SubscriptionMode subscriptionMode, SubscriptionType subscriptionType) {
+    private static Subscription generateSubscription(String topic, 
SubscriptionMode subscriptionMode,
+                                                     SubscriptionType 
subscriptionType) {
         Subscription subscription = new Subscription();
         List<SubscriptionItem> subscriptionItems = new ArrayList<>();
         subscriptionItems.add(new SubscriptionItem(topic, subscriptionMode, 
subscriptionType));
diff --git 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
index 1f4e774..baa5013 100644
--- 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
+++ 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
@@ -9,9 +9,9 @@ import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
 import org.apache.eventmesh.client.tcp.common.RequestContext;
 import org.apache.eventmesh.client.tcp.common.TcpClient;
 import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
-import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.exception.EventMeshException;
 import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Package;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 
@@ -82,6 +82,7 @@ public class EventMeshMessageTCPPubClient extends TcpClient 
implements EventMesh
         }
     }
 
+    // todo: Maybe use org.apache.eventmesh.common.EvetMesh here is better
     @Override
     public Package rr(EventMeshMessage eventMeshMessage, long timeout) throws 
EventMeshException {
         try {
diff --git 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
index cdce845..5891e4b 100644
--- 
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
+++ 
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
@@ -7,12 +7,12 @@ import org.apache.eventmesh.client.tcp.common.ReceiveMsgHook;
 import org.apache.eventmesh.client.tcp.common.RequestContext;
 import org.apache.eventmesh.client.tcp.common.TcpClient;
 import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
-import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.exception.EventMeshException;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
 import org.apache.eventmesh.common.protocol.tcp.Command;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Package;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
index 0b26161..30e7142 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
@@ -17,14 +17,15 @@
 
 package org.apache.eventmesh.client.http.demo;
 
-import org.apache.commons.lang3.StringUtils;
-
 import org.apache.eventmesh.client.http.conf.EventMeshHttpClientConfig;
 import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
-import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.EventMeshMessage;
+import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
 import org.apache.eventmesh.common.utils.ThreadUtils;
+
+import org.apache.commons.lang3.StringUtils;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -45,38 +46,36 @@ public class SyncRequestInstance {
                 eventMeshIPPort = "127.0.0.1:10105";
             }
 
-            EventMeshHttpClientConfig eventMeshClientConfig = new 
EventMeshHttpClientConfig();
-            eventMeshClientConfig.setLiteEventMeshAddr(eventMeshIPPort)
-                    .setProducerGroup("EventMeshTest-producerGroup")
-                    .setEnv("env")
-                    .setIdc("idc")
-                    .setIp(IPUtils.getLocalAddress())
-                    .setSys("1234")
-                    .setPid(String.valueOf(ThreadUtils.getPID()));
+            EventMeshHttpClientConfig eventMeshClientConfig = 
EventMeshHttpClientConfig.builder()
+                .liteEventMeshAddr(eventMeshIPPort)
+                .producerGroup("EventMeshTest-producerGroup")
+                .env("env")
+                .idc("idc")
+                .ip(IPUtils.getLocalAddress())
+                .sys("1234")
+                .pid(String.valueOf(ThreadUtils.getPID())).build();
 
             eventMeshHttpProducer = new 
EventMeshHttpProducer(eventMeshClientConfig);
-            eventMeshHttpProducer.start();
 
             long startTime = System.currentTimeMillis();
-            EventMeshMessage eventMeshMessage = new EventMeshMessage();
-            eventMeshMessage.setBizSeqNo(RandomStringUtils.generateNum(30))
-                    .setContent("contentStr with special protocal")
-                    .setTopic(topic)
-                    .setUniqueId(RandomStringUtils.generateNum(30));
+            EventMeshMessage eventMeshMessage = EventMeshMessage.builder()
+                .bizSeqNo(RandomStringUtils.generateNum(30))
+                .content("contentStr with special protocal")
+                .topic(topic)
+                .uniqueId(RandomStringUtils.generateNum(30)).build();
 
             EventMeshMessage rsp = 
eventMeshHttpProducer.request(eventMeshMessage, 10000);
             if (logger.isDebugEnabled()) {
-                logger.debug("sendmsg : {}, return : {}, cost:{}ms", 
eventMeshMessage.getContent(), rsp.getContent(), System.currentTimeMillis() - 
startTime);
+                logger.debug("sendmsg : {}, return : {}, cost:{}ms", 
eventMeshMessage.getContent(), rsp.getContent(),
+                    System.currentTimeMillis() - startTime);
             }
         } catch (Exception e) {
             logger.warn("send msg failed", e);
         }
 
-        try {
-            Thread.sleep(30000);
-            if (eventMeshHttpProducer != null) {
-                eventMeshHttpProducer.shutdown();
-            }
+        Thread.sleep(30000);
+        try (final EventMeshHttpProducer closed = eventMeshHttpProducer){
+            // close producer
         } catch (Exception e1) {
             logger.warn("producer shutdown exception", e1);
         }
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
index 4c03ee1..aabafb1 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/util/HttpLoadBalanceUtilsTest.java
@@ -28,8 +28,9 @@ public class HttpLoadBalanceUtilsTest {
 
     @Test
     public void testCreateRandomSelector() throws EventMeshException {
-        EventMeshHttpClientConfig eventMeshHttpClientConfig = new 
EventMeshHttpClientConfig()
-                .setLiteEventMeshAddr("127.0.0.1:1001;127.0.0.2:1002");
+        EventMeshHttpClientConfig eventMeshHttpClientConfig = 
EventMeshHttpClientConfig.builder()
+            .liteEventMeshAddr("127.0.0.1:1001;127.0.0.2:1002")
+            .build();
         LoadBalanceSelector<String> randomSelector = HttpLoadBalanceUtils
                 
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
         Assert.assertEquals(LoadBalanceType.RANDOM, randomSelector.getType());
@@ -37,9 +38,9 @@ public class HttpLoadBalanceUtilsTest {
 
     @Test
     public void testCreateWeightRoundRobinSelector() throws EventMeshException 
{
-        EventMeshHttpClientConfig eventMeshHttpClientConfig = new 
EventMeshHttpClientConfig()
-                .setLiteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
-                .setLoadBalanceType(LoadBalanceType.WEIGHT_ROUND_ROBIN);
+        EventMeshHttpClientConfig eventMeshHttpClientConfig = 
EventMeshHttpClientConfig.builder()
+                .liteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
+                .loadBalanceType(LoadBalanceType.WEIGHT_ROUND_ROBIN).build();
         LoadBalanceSelector<String> weightRoundRobinSelector = 
HttpLoadBalanceUtils
                 
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
         Assert.assertEquals(LoadBalanceType.WEIGHT_ROUND_ROBIN, 
weightRoundRobinSelector.getType());
@@ -47,9 +48,9 @@ public class HttpLoadBalanceUtilsTest {
 
     @Test
     public void testCreateWeightRandomSelector() throws EventMeshException {
-        EventMeshHttpClientConfig eventMeshHttpClientConfig = new 
EventMeshHttpClientConfig()
-                .setLiteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
-                .setLoadBalanceType(LoadBalanceType.WEIGHT_RANDOM);
+        EventMeshHttpClientConfig eventMeshHttpClientConfig = 
EventMeshHttpClientConfig.builder()
+                .liteEventMeshAddr("127.0.0.1:1001:1;127.0.0.2:1001:2")
+                .loadBalanceType(LoadBalanceType.WEIGHT_RANDOM).build();
         LoadBalanceSelector<String> weightRoundRobinSelector = 
HttpLoadBalanceUtils
                 
.createEventMeshServerLoadBalanceSelector(eventMeshHttpClientConfig);
         Assert.assertEquals(LoadBalanceType.WEIGHT_RANDOM, 
weightRoundRobinSelector.getType());
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
index 78e1016..f6397b8 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/common/EventMeshTestUtils.java
@@ -20,53 +20,48 @@ package org.apache.eventmesh.client.tcp.common;
 import static 
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_SyncSubscribeTest;
 import static 
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientBroadCast;
 import static 
org.apache.eventmesh.client.tcp.common.EventMeshTestCaseTopicSet.TOPIC_PRX_WQ2ClientUniCast;
-
-import java.util.concurrent.ThreadLocalRandom;
+import static 
org.apache.eventmesh.common.protocol.tcp.Command.RESPONSE_TO_SERVER;
 
 import org.apache.eventmesh.common.protocol.tcp.Command;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Header;
-import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
 
-import static 
org.apache.eventmesh.common.protocol.tcp.Command.RESPONSE_TO_SERVER;
-
-
-
+import java.util.concurrent.ThreadLocalRandom;
 
 public class EventMeshTestUtils {
     private static final int seqLength = 10;
 
     public static UserAgent generateClient1() {
-        UserAgent user = new UserAgent();
-        user.setHost("127.0.0.1");
-        user.setPassword(generateRandomString(8));
-        user.setUsername("PU4283");
-        user.setConsumerGroup("EventmeshTest-ConsumerGroup");
-        user.setProducerGroup("EventmeshTest-ProducerGroup");
-        user.setPath("/data/app/umg_proxy");
-        user.setPort(8362);
-        user.setSubsystem("5023");
-        user.setPid(32893);
-        user.setVersion("2.0.11");
-        user.setIdc("FT");
-        return user;
+        return UserAgent.builder()
+            .host("127.0.0.1")
+            .password(generateRandomString(8))
+            .username("PU4283")
+            .consumerGroup("EventmeshTest-ConsumerGroup")
+            .producerGroup("EventmeshTest-ProducerGroup")
+            .path("/data/app/umg_proxy")
+            .port(8362)
+            .subsystem("5023")
+            .pid(32893)
+            .version("2.0.11")
+            .idc("FT")
+            .build();
     }
 
     public static UserAgent generateClient2() {
-        UserAgent user = new UserAgent();
-        user.setHost("127.0.0.1");
-        user.setPassword(generateRandomString(8));
-        user.setUsername("PU4283");
-        user.setConsumerGroup("EventmeshTest-ConsumerGroup");
-        user.setProducerGroup("EventmeshTest-ProducerGroup");
-        user.setPath("/data/app/umg_proxy");
-        user.setPort(9362);
-        user.setSubsystem("5017");
-        user.setPid(42893);
-        user.setVersion("2.0.11");
-        user.setIdc("FT");
-        return user;
+        return UserAgent.builder()
+            .host("127.0.0.1")
+            .password(generateRandomString(8))
+            .username("PU4283")
+            .consumerGroup("EventmeshTest-ConsumerGroup")
+            .producerGroup("EventmeshTest-ProducerGroup")
+            .path("/data/app/umg_proxy")
+            .port(9362)
+            .subsystem("5017")
+            .pid(42893)
+            .version("2.0.11")
+            .idc("FT").build();
     }
 
     public static Package syncRR() {
@@ -104,7 +99,7 @@ public class EventMeshTestUtils {
         return msg;
     }
 
-    private static EventMeshMessage generateSyncRRMqMsg() {
+    public static EventMeshMessage generateSyncRRMqMsg() {
         EventMeshMessage mqMsg = new EventMeshMessage();
         mqMsg.setTopic(TOPIC_PRX_SyncSubscribeTest);
         mqMsg.getProperties().put("msgType", "persistent");
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
index 272ece0..7b6b4ca 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/tcp/demo/SyncRequest.java
@@ -17,12 +17,14 @@
 
 package org.apache.eventmesh.client.tcp.demo;
 
-import org.apache.eventmesh.client.tcp.EventMeshTCPClient;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
 import org.apache.eventmesh.client.tcp.common.EventMeshTestUtils;
-import org.apache.eventmesh.client.tcp.impl.DefaultEventMeshTCPClient;
-import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import org.apache.eventmesh.client.tcp.conf.EventMeshTcpClientConfig;
+import 
org.apache.eventmesh.client.tcp.impl.eventmeshmessage.EventMeshMessageTCPPubClient;
+import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.Package;
+import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -30,18 +32,23 @@ public class SyncRequest {
 
     public static Logger logger = LoggerFactory.getLogger(SyncRequest.class);
 
-    private static EventMeshTCPClient client;
+    private static EventMeshMessageTCPPubClient client;
 
-    public static void main(String[] agrs) throws Exception {
+    public static void main(String[] agrs) {
         try {
             UserAgent userAgent = EventMeshTestUtils.generateClient1();
-            client = new DefaultEventMeshTCPClient("127.0.0.1", 10000, 
userAgent);
+            EventMeshTcpClientConfig eventMeshTcpClientConfig = 
EventMeshTcpClientConfig.builder()
+                .host("127.0.0.1")
+                .port(10000)
+                .userAgent(userAgent)
+                .build();
+            client = new 
EventMeshMessageTCPPubClient(eventMeshTcpClientConfig);
             client.init();
             client.heartbeat();
 
-            Package rrMsg = EventMeshTestUtils.syncRR();
-            logger.info("begin send rr msg=================={}", rrMsg);
-            Package response = client.rr(rrMsg, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
+            EventMeshMessage eventMeshMessage = 
EventMeshTestUtils.generateSyncRRMqMsg();
+            logger.info("begin send rr msg=================={}", 
eventMeshMessage);
+            Package response = client.rr(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
             logger.info("receive rr reply==================={}", response);
 
             // release resource and close client

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

Reply via email to