This is an automated email from the ASF dual-hosted git repository.
mikexue pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
The following commit(s) were added to refs/heads/develop by this push:
new 0acafcc resolve_exception->Consumer subscription topic is invalid
#590 (#598)
0acafcc is described below
commit 0acafccc080e0d31d7d24e6384ca4527702354a0
Author: hagsyn <[email protected]>
AuthorDate: Fri Nov 19 17:22:34 2021 +0800
resolve_exception->Consumer subscription topic is invalid #590 (#598)
close #590
---
.../http/processor/SubscribeProcessor.java | 115 +++++++++++----------
1 file changed, 58 insertions(+), 57 deletions(-)
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
index d91b462..ab411f9 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/processor/SubscribeProcessor.java
@@ -70,49 +70,49 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
@Override
public void processRequest(ChannelHandlerContext ctx,
AsyncContext<HttpCommand> asyncContext)
- throws Exception {
+ throws Exception {
HttpCommand responseEventMeshCommand;
httpLogger.info("cmd={}|{}|client2eventMesh|from={}|to={}",
-
RequestCode.get(Integer.valueOf(asyncContext.getRequest().getRequestCode())),
- EventMeshConstants.PROTOCOL_HTTP,
- RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
IPUtil.getLocalAddress());
+
RequestCode.get(Integer.valueOf(asyncContext.getRequest().getRequestCode())),
+ EventMeshConstants.PROTOCOL_HTTP,
+ RemotingHelper.parseChannelRemoteAddr(ctx.channel()),
IPUtil.getLocalAddress());
SubscribeRequestHeader subscribeRequestHeader =
- (SubscribeRequestHeader) asyncContext.getRequest().getHeader();
+ (SubscribeRequestHeader) asyncContext.getRequest().getHeader();
SubscribeRequestBody subscribeRequestBody =
- (SubscribeRequestBody) asyncContext.getRequest().getBody();
+ (SubscribeRequestBody) asyncContext.getRequest().getBody();
SubscribeResponseHeader subscribeResponseHeader =
- SubscribeResponseHeader
-
.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
-
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
- IPUtil.getLocalAddress(),
-
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
-
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
+ SubscribeResponseHeader
+
.buildHeader(Integer.valueOf(asyncContext.getRequest().getRequestCode()),
+
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshCluster,
+ IPUtil.getLocalAddress(),
+
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshEnv,
+
eventMeshHTTPServer.getEventMeshHttpConfiguration().eventMeshIDC);
//validate header
if (StringUtils.isBlank(subscribeRequestHeader.getIdc())
- || StringUtils.isBlank(subscribeRequestHeader.getPid())
- || !StringUtils.isNumeric(subscribeRequestHeader.getPid())
- || StringUtils.isBlank(subscribeRequestHeader.getSys())) {
+ || StringUtils.isBlank(subscribeRequestHeader.getPid())
+ || !StringUtils.isNumeric(subscribeRequestHeader.getPid())
+ || StringUtils.isBlank(subscribeRequestHeader.getSys())) {
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
- subscribeResponseHeader,
- SubscribeResponseBody
-
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getRetCode(),
-
EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getErrMsg()));
+ subscribeResponseHeader,
+ SubscribeResponseBody
+
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getRetCode(),
+
EventMeshRetCode.EVENTMESH_PROTOCOL_HEADER_ERR.getErrMsg()));
asyncContext.onComplete(responseEventMeshCommand);
return;
}
//validate body
if (StringUtils.isBlank(subscribeRequestBody.getUrl())
- || CollectionUtils.isEmpty(subscribeRequestBody.getTopics())
- || StringUtils.isBlank(subscribeRequestBody.getConsumerGroup())) {
+ || CollectionUtils.isEmpty(subscribeRequestBody.getTopics())
+ ||
StringUtils.isBlank(subscribeRequestBody.getConsumerGroup())) {
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
- subscribeResponseHeader,
- SubscribeResponseBody
-
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getRetCode(),
-
EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getErrMsg()));
+ subscribeResponseHeader,
+ SubscribeResponseBody
+
.buildBody(EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getRetCode(),
+
EventMeshRetCode.EVENTMESH_PROTOCOL_BODY_ERR.getErrMsg()));
asyncContext.onComplete(responseEventMeshCommand);
return;
}
@@ -128,17 +128,17 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
for (SubscriptionItem item : subTopicList) {
try {
Acl.doAclCheckInHttpReceive(remoteAddr, user, pass,
subsystem, item.getTopic(),
- requestCode);
+ requestCode);
} catch (Exception e) {
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
- subscribeResponseHeader,
- SendMessageResponseBody
-
.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
- e.getMessage()));
+ subscribeResponseHeader,
+ SendMessageResponseBody
+
.buildBody(EventMeshRetCode.EVENTMESH_ACL_ERR.getRetCode(),
+ e.getMessage()));
asyncContext.onComplete(responseEventMeshCommand);
aclLogger
- .warn("CLIENT HAS NO PERMISSION,SubscribeProcessor
subscribe failed", e);
+ .warn("CLIENT HAS NO PERMISSION,SubscribeProcessor
subscribe failed", e);
return;
}
}
@@ -153,7 +153,7 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
for (SubscriptionItem subTopic : subTopicList) {
List<Client> groupTopicClients =
eventMeshHTTPServer.localClientInfoMapping
- .get(consumerGroup + "@" + subTopic.getTopic());
+ .get(consumerGroup + "@" + subTopic.getTopic());
if (CollectionUtils.isEmpty(groupTopicClients)) {
httpLogger.error("group {} topic {} clients is empty",
consumerGroup, subTopic);
@@ -170,7 +170,7 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
}
}
ConsumerGroupConf consumerGroupConf =
-
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup);
+
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup);
if (consumerGroupConf == null) {
// new subscription
consumerGroupConf = new ConsumerGroupConf(consumerGroup);
@@ -188,7 +188,17 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
} else {
// already subscribed
Map<String, ConsumerGroupTopicConf> map =
- consumerGroupConf.getConsumerGroupTopicConf();
+ consumerGroupConf.getConsumerGroupTopicConf();
+ if (!map.containsKey(subTopic.getTopic())) {
+ //If there are multiple topics, append it
+ ConsumerGroupTopicConf newTopicConf = new
ConsumerGroupTopicConf();
+ newTopicConf.setConsumerGroup(consumerGroup);
+ newTopicConf.setTopic(subTopic.getTopic());
+ newTopicConf.setSubscriptionItem(subTopic);
+ newTopicConf.setUrls(new
HashSet<>(Arrays.asList(url)));
+ newTopicConf.setIdcUrls(idcUrls);
+ map.put(subTopic.getTopic(), newTopicConf);
+ }
for (String key : map.keySet()) {
if (StringUtils.equals(subTopic.getTopic(), key)) {
ConsumerGroupTopicConf latestTopicConf = new
ConsumerGroupTopicConf();
@@ -202,15 +212,6 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
latestTopicConf.setIdcUrls(idcUrls);
map.put(key, latestTopicConf);
- } else {
- //If there are multiple topics, append it
- ConsumerGroupTopicConf newTopicConf = new
ConsumerGroupTopicConf();
- newTopicConf.setConsumerGroup(consumerGroup);
- newTopicConf.setTopic(subTopic.getTopic());
- newTopicConf.setSubscriptionItem(subTopic);
- newTopicConf.setUrls(new
HashSet<>(Arrays.asList(url)));
- newTopicConf.setIdcUrls(idcUrls);
- map.put(subTopic.getTopic(), newTopicConf);
}
}
}
@@ -221,7 +222,7 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
try {
// subscription relationship change notification
eventMeshHTTPServer.getConsumerManager().notifyConsumerManager(consumerGroup,
-
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup));
+
eventMeshHTTPServer.localConsumerGroupMapping.get(consumerGroup));
final CompleteHandler<HttpCommand> handler = new
CompleteHandler<HttpCommand>() {
@Override
@@ -232,8 +233,8 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
}
eventMeshHTTPServer.sendResponse(ctx,
httpCommand.httpResponse());
eventMeshHTTPServer.metrics.summaryMetrics.recordHTTPReqResTimeCost(
- System.currentTimeMillis()
- - asyncContext.getRequest().getReqTime());
+ System.currentTimeMillis()
+ -
asyncContext.getRequest().getReqTime());
} catch (Exception ex) {
// ignore
}
@@ -241,22 +242,22 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
};
responseEventMeshCommand =
asyncContext.getRequest().createHttpCommandResponse(
- EventMeshRetCode.SUCCESS.getRetCode(),
EventMeshRetCode.SUCCESS.getErrMsg());
+ EventMeshRetCode.SUCCESS.getRetCode(),
EventMeshRetCode.SUCCESS.getErrMsg());
asyncContext.onComplete(responseEventMeshCommand, handler);
} catch (Exception e) {
HttpCommand err =
asyncContext.getRequest().createHttpCommandResponse(
- subscribeResponseHeader,
- SubscribeResponseBody
-
.buildBody(EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
-
EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getErrMsg()
- + EventMeshUtil.stackTrace(e, 2)));
+ subscribeResponseHeader,
+ SubscribeResponseBody
+
.buildBody(EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getRetCode(),
+
EventMeshRetCode.EVENTMESH_SUBSCRIBE_ERR.getErrMsg()
+ + EventMeshUtil.stackTrace(e,
2)));
asyncContext.onComplete(err);
long endTime = System.currentTimeMillis();
httpLogger.error(
- "message|eventMesh2mq|REQ|ASYNC|send2MQCost={}ms|topic={}"
- + "|bizSeqNo={}|uniqueId={}", endTime - startTime,
- JsonUtils.serialize(subscribeRequestBody.getTopics()),
- subscribeRequestBody.getUrl(), e);
+
"message|eventMesh2mq|REQ|ASYNC|send2MQCost={}ms|topic={}"
+ + "|bizSeqNo={}|uniqueId={}", endTime -
startTime,
+ JsonUtils.serialize(subscribeRequestBody.getTopics()),
+ subscribeRequestBody.getUrl(), e);
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgFailed();
eventMeshHTTPServer.metrics.summaryMetrics.recordSendMsgCost(endTime -
startTime);
}
@@ -286,7 +287,7 @@ public class SubscribeProcessor implements
HttpRequestProcessor {
if
(eventMeshHTTPServer.localClientInfoMapping.containsKey(groupTopicKey)) {
List<Client> localClients =
-
eventMeshHTTPServer.localClientInfoMapping.get(groupTopicKey);
+
eventMeshHTTPServer.localClientInfoMapping.get(groupTopicKey);
boolean isContains = false;
for (Client localClient : localClients) {
if (StringUtils.equals(localClient.url, client.url)) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]