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