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 deee538 [Feature #562] support cloudevents api in
eventmesh-connector-api module (#578)
deee538 is described below
commit deee538f5ed0f65f16f387de42e68421b2dd7f68
Author: mike_xwm <[email protected]>
AuthorDate: Wed Nov 10 10:23:32 2021 +0800
[Feature #562] support cloudevents api in eventmesh-connector-api module
(#578)
* [Feature #564] Support CloudEvents protocols for pub/sub in
EventMesh-feature design
* support cloudevents api in eventmesh-connector-api module
* fix checkStyle
* fix checkStyle
* fix checkStyle
* 1.support LifeCycle.java
2.update Consumer and Producer
* fix remove the extra blank line
---
.../eventmesh-connector-api/build.gradle | 1 +
...entMeshAction.java => AsyncConsumeContext.java} | 8 ++--
.../{EventMeshAction.java => EventListener.java} | 20 ++++++--
.../org/apache/eventmesh/api/EventMeshAction.java | 1 +
.../api/EventMeshAsyncConsumeContext.java | 1 +
.../api/{EventMeshAction.java => LifeCycle.java} | 19 ++++++--
.../{EventMeshAction.java => SendCallback.java} | 16 +++++--
.../api/{EventMeshAction.java => SendResult.java} | 28 +++++++++--
.../Consumer.java} | 40 +++++++++-------
.../ConnectorRuntimeException.java} | 23 +++++++--
.../OnExceptionContext.java} | 44 ++++++++++++------
.../apache/eventmesh/api/producer/Producer.java | 54 ++++++++++++++++++++++
12 files changed, 199 insertions(+), 56 deletions(-)
diff --git a/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
b/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
index 19cb54c..a721189 100644
--- a/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
+++ b/eventmesh-connector-plugin/eventmesh-connector-api/build.gradle
@@ -18,6 +18,7 @@
dependencies {
implementation project(":eventmesh-spi")
implementation project(":eventmesh-common")
+ api 'io.cloudevents:cloudevents-core'
api 'io.openmessaging:openmessaging-api'
api 'io.dropwizard.metrics:metrics-core'
api "io.dropwizard.metrics:metrics-healthchecks"
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
similarity index 89%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
index 4fda6d0..7c1c739 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/AsyncConsumeContext.java
@@ -14,12 +14,12 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
- ReconsumeLater,
+public abstract class AsyncConsumeContext {
+
+ public abstract void commit(EventMeshAction action);
- ManualAck
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
similarity index 67%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
index 4fda6d0..eede41c 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventListener.java
@@ -14,12 +14,24 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
- ReconsumeLater,
+import io.cloudevents.CloudEvent;
+
+/**
+ * Event listener, registered for consume messages by consumer.
+ *
+ * <p>
+ * <strong>
+ * Thread safe requirements: this interface will be invoked by multi threads,
+ * so users should keep thread safe during the consume process.
+ * </strong>
+ * </p>
+ */
+public interface EventListener {
+
+ void consume(final CloudEvent cloudEvent, final AsyncConsumeContext
context);
- ManualAck
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
index 4fda6d0..f783201 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
@@ -14,6 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
public enum EventMeshAction {
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
index c7e4e7f..c123ca2 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
@@ -14,6 +14,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
import io.openmessaging.api.Action;
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
similarity index 69%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
index 4fda6d0..f5d61de 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/LifeCycle.java
@@ -14,12 +14,23 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
+import org.apache.eventmesh.api.consumer.Consumer;
+import org.apache.eventmesh.api.producer.Producer;
+
+/**
+ * The {@code LifeCycle} defines a lifecycle interface for a OMS related
service endpoint,
+ * like {@link Producer}, {@link Consumer}, and so on.
+ */
+public interface LifeCycle {
+
+ boolean isStarted();
+
+ boolean isClosed();
- ReconsumeLater,
+ void start();
- ManualAck
+ void shutdown();
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
similarity index 68%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
index 4fda6d0..c955f5d 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendCallback.java
@@ -14,12 +14,20 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
+import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.api.producer.Producer;
+
+import io.cloudevents.CloudEvent;
+
+/**
+ * Call back interface used in {@link Producer#sendAsync(CloudEvent,
SendCallback)}.
+ */
+public interface SendCallback {
- ReconsumeLater,
+ void onSuccess(final SendResult sendResult);
- ManualAck
+ void onException(final OnExceptionContext context);
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
similarity index 62%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
index 4fda6d0..1a68132 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/SendResult.java
@@ -14,12 +14,32 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
+
package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
+public class SendResult {
+ private String messageId;
+
+ private String topic;
+
+ public String getMessageId() {
+ return messageId;
+ }
+
+ public void setMessageId(String messageId) {
+ this.messageId = messageId;
+ }
+
+ public String getTopic() {
+ return topic;
+ }
- ReconsumeLater,
+ public void setTopic(String topic) {
+ this.topic = topic;
+ }
- ManualAck
+ @Override
+ public String toString() {
+ return "SendResult[topic=" + topic + ", messageId=" + messageId + ']';
+ }
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
similarity index 50%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
index c7e4e7f..a6fafe0 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/consumer/Consumer.java
@@ -14,27 +14,33 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.eventmesh.api;
-import io.openmessaging.api.Action;
-import io.openmessaging.api.AsyncConsumeContext;
+package org.apache.eventmesh.api.consumer;
-public abstract class EventMeshAsyncConsumeContext extends AsyncConsumeContext
{
+import org.apache.eventmesh.api.AbstractContext;
+import org.apache.eventmesh.api.EventListener;
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.spi.EventMeshExtensionType;
+import org.apache.eventmesh.spi.EventMeshSPI;
- private AbstractContext abstractContext;
+import java.util.List;
+import java.util.Properties;
- public AbstractContext getAbstractContext() {
- return abstractContext;
- }
+import io.cloudevents.CloudEvent;
- public void setAbstractContext(AbstractContext abstractContext) {
- this.abstractContext = abstractContext;
- }
- public abstract void commit(EventMeshAction action);
- @Override
- public void commit(Action action) {
- throw new UnsupportedOperationException("not support yet");
- }
-}
\ No newline at end of file
+/**
+ * Consumer Interface.
+ */
+@EventMeshSPI(isSingleton = false, eventMeshExtensionType =
EventMeshExtensionType.CONNECTOR)
+public interface Consumer extends LifeCycle {
+
+ void init(Properties keyValue) throws Exception;
+
+ void updateOffset(List<CloudEvent> cloudEvents, AbstractContext context);
+
+ void subscribe(String topic, final EventListener listener) throws
Exception;
+
+ void unsubscribe(String topic);
+}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
similarity index 63%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
index 4fda6d0..0efe513 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAction.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/ConnectorRuntimeException.java
@@ -14,12 +14,25 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.eventmesh.api;
-public enum EventMeshAction {
- CommitMessage,
+package org.apache.eventmesh.api.exception;
- ReconsumeLater,
+public class ConnectorRuntimeException extends RuntimeException {
+
+ public ConnectorRuntimeException() {
+
+ }
+
+ public ConnectorRuntimeException(String message) {
+ super(message);
+ }
+
+ public ConnectorRuntimeException(Throwable throwable) {
+ super(throwable);
+ }
+
+ public ConnectorRuntimeException(String message, Throwable throwable) {
+ super(message, throwable);
+ }
- ManualAck
}
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
similarity index 53%
copy from
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
copy to
eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
index c7e4e7f..1237e5f 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/EventMeshAsyncConsumeContext.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/exception/OnExceptionContext.java
@@ -14,27 +14,43 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.eventmesh.api;
-import io.openmessaging.api.Action;
-import io.openmessaging.api.AsyncConsumeContext;
+package org.apache.eventmesh.api.exception;
-public abstract class EventMeshAsyncConsumeContext extends AsyncConsumeContext
{
+public class OnExceptionContext {
- private AbstractContext abstractContext;
+ private String messageId;
- public AbstractContext getAbstractContext() {
- return abstractContext;
+ private String topic;
+
+ /**
+ * Detailed exception stack information.
+ */
+ private ConnectorRuntimeException exception;
+
+
+ public String getMessageId() {
+ return messageId;
+ }
+
+ public void setMessageId(String messageId) {
+ this.messageId = messageId;
}
- public void setAbstractContext(AbstractContext abstractContext) {
- this.abstractContext = abstractContext;
+ public String getTopic() {
+ return topic;
}
- public abstract void commit(EventMeshAction action);
+ public void setTopic(String topic) {
+ this.topic = topic;
+ }
+
+
+ public ConnectorRuntimeException getException() {
+ return exception;
+ }
- @Override
- public void commit(Action action) {
- throw new UnsupportedOperationException("not support yet");
+ public void setException(ConnectorRuntimeException exception) {
+ this.exception = exception;
}
-}
\ No newline at end of file
+}
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
new file mode 100644
index 0000000..08d8998
--- /dev/null
+++
b/eventmesh-connector-plugin/eventmesh-connector-api/src/main/java/org/apache/eventmesh/api/producer/Producer.java
@@ -0,0 +1,54 @@
+/*
+ * 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.api.producer;
+
+import org.apache.eventmesh.api.LifeCycle;
+import org.apache.eventmesh.api.RRCallback;
+import org.apache.eventmesh.api.SendCallback;
+import org.apache.eventmesh.api.SendResult;
+import org.apache.eventmesh.spi.EventMeshExtensionType;
+import org.apache.eventmesh.spi.EventMeshSPI;
+
+import java.util.Properties;
+
+import io.cloudevents.CloudEvent;
+
+/**
+ * Producer Interface.
+ */
+@EventMeshSPI(isSingleton = false, eventMeshExtensionType =
EventMeshExtensionType.CONNECTOR)
+public interface Producer extends LifeCycle {
+
+ void init(Properties properties) throws Exception;
+
+ SendResult publish(final CloudEvent cloudEvent);
+
+ void publish(CloudEvent cloudEvent, SendCallback sendCallback) throws
Exception;
+
+ void sendOneway(final CloudEvent cloudEvent);
+
+ void sendAsync(final CloudEvent cloudEvent, final SendCallback
sendCallback);
+
+ void request(CloudEvent cloudEvent, RRCallback rrCallback, long timeout)
throws Exception;
+
+ boolean reply(final CloudEvent cloudEvent, final SendCallback
sendCallback) throws Exception;
+
+ void checkTopicExist(String topic) throws Exception;
+
+ void setExtFields();
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]