This is an automated email from the ASF dual-hosted git repository.
jonyang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git
The following commit(s) were added to refs/heads/master by this push:
new 47f4d9fb8 Simplify code with lambda (#3130)
47f4d9fb8 is described below
commit 47f4d9fb8dc0828bd9764a58ae8623f188e07dbd
Author: weihubeats <[email protected]>
AuthorDate: Mon Feb 13 14:08:57 2023 +0800
Simplify code with lambda (#3130)
---
.../ConsumeMessageConcurrentlyService.java | 74 +++++++++-------------
1 file changed, 29 insertions(+), 45 deletions(-)
diff --git
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
index e2119ec2e..0b00b4fe2 100644
---
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
+++
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
@@ -43,6 +43,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
+import java.util.Objects;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.Executors;
import java.util.concurrent.LinkedBlockingQueue;
@@ -74,7 +75,7 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
this.defaultMQPushConsumer =
this.defaultMQPushConsumerImpl.getDefaultMQPushConsumer();
this.consumerGroup = this.defaultMQPushConsumer.getConsumerGroup();
- this.consumeRequestQueue = new LinkedBlockingQueue<Runnable>();
+ this.consumeRequestQueue = new LinkedBlockingQueue<>();
this.consumeExecutor = new ThreadPoolExecutor(
this.defaultMQPushConsumer.getConsumeThreadMin(),
@@ -92,17 +93,12 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
public void start() {
if (this.defaultMQPushConsumer.getConsumeTimeout() > 0) {
- this.cleanExpireMsgExecutors.scheduleAtFixedRate(new Runnable() {
-
- @Override
- public void run() {
- try {
- cleanExpireMsg();
- } catch (Exception e) {
- log.warn("cleanExpireMsg ", e);
- }
+ this.cleanExpireMsgExecutors.scheduleAtFixedRate(() -> {
+ try {
+ cleanExpireMsg();
+ } catch (Exception e) {
+ log.warn("cleanExpireMsg ", e);
}
-
}, this.defaultMQPushConsumer.getConsumeTimeout(),
this.defaultMQPushConsumer.getConsumeTimeout(), TimeUnit.MINUTES);
}
}
@@ -215,7 +211,7 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
}
} else {
for (int total = 0; total < msgs.size(); ) {
- List<MessageExt> msgThis = new
ArrayList<MessageExt>(consumeBatchSize);
+ List<MessageExt> msgThis = new ArrayList<>(consumeBatchSize);
for (int i = 0; i < consumeBatchSize; i++, total++) {
if (total < msgs.size()) {
msgThis.add(msgs.get(total));
@@ -352,40 +348,31 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
final MessageQueue messageQueue
) {
- this.scheduledExecutorService.schedule(new Runnable() {
-
- @Override
- public void run() {
-
ConsumeMessageConcurrentlyService.this.submitConsumeRequest(msgs, processQueue,
messageQueue, true);
- }
- }, 5000, TimeUnit.MILLISECONDS);
+ this.scheduledExecutorService.schedule(
+ () ->
ConsumeMessageConcurrentlyService.this.submitConsumeRequest(msgs, processQueue,
messageQueue, true), 5000, TimeUnit.MILLISECONDS);
}
private void submitConsumeRequestLater(final ConsumeRequest
consumeRequest) {
final int times =
defaultMQPushConsumerImpl.getDefaultMQPushConsumer().getMaxReconsumeTimes();
log.warn("rejected by thread pool, try resubmit {} times,
consumerGroup:{}", times,
defaultMQPushConsumerImpl.getDefaultMQPushConsumer().getConsumerGroup());
- this.scheduledExecutorService.schedule(new Runnable() {
-
- @Override
- public void run() {
- boolean success = false;
- for (int i = 0; i < times; i++) {
- try {
- ThreadUtils.sleep(1, TimeUnit.SECONDS);
-
ConsumeMessageConcurrentlyService.this.consumeExecutor.submit(consumeRequest);
- success = true;
- break;
- } catch (RejectedExecutionException e) {
- //ignore
- }
+ this.scheduledExecutorService.schedule(() -> {
+ boolean success = false;
+ for (int i = 0; i < times; i++) {
+ try {
+ ThreadUtils.sleep(1, TimeUnit.SECONDS);
+
ConsumeMessageConcurrentlyService.this.consumeExecutor.submit(consumeRequest);
+ success = true;
+ break;
+ } catch (RejectedExecutionException e) {
+ //ignore
}
- if (!success) {
- for (MessageExt messageExt : consumeRequest.getMsgs()) {
- log.warn("discard rejected messages {} after retry {}
times", messageExt, times);
- }
-
consumeRequest.getProcessQueue().removeMessage(consumeRequest.getMsgs());
+ }
+ if (!success) {
+ for (MessageExt messageExt : consumeRequest.getMsgs()) {
+ log.warn("discard rejected messages {} after retry {}
times", messageExt, times);
}
+
consumeRequest.getProcessQueue().removeMessage(consumeRequest.getMsgs());
}
}, 1000, TimeUnit.MILLISECONDS);
}
@@ -417,7 +404,6 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
return;
}
- MessageListenerConcurrently listener =
ConsumeMessageConcurrentlyService.this.messageListener;
EventMeshConsumeConcurrentlyContext context = new
EventMeshConsumeConcurrentlyContext(messageQueue, processQueue);
ConsumeConcurrentlyStatus status = null;
@@ -425,7 +411,7 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
if
(ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.hasHook()) {
consumeMessageContext = new ConsumeMessageContext();
consumeMessageContext.setConsumerGroup(defaultMQPushConsumer.getConsumerGroup());
- consumeMessageContext.setProps(new HashMap<String, String>());
+ consumeMessageContext.setProps(new HashMap<>());
consumeMessageContext.setMq(messageQueue);
consumeMessageContext.setMsgList(msgs);
consumeMessageContext.setSuccess(false);
@@ -444,7 +430,7 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
}
- status =
listener.consumeMessage(Collections.unmodifiableList(msgs), context);
+ status =
ConsumeMessageConcurrentlyService.this.messageListener.consumeMessage(Collections.unmodifiableList(msgs),
context);
} catch (Throwable e) {
log.warn("consumeMessage exception: {} Group: {} Msgs: {} MQ:
{}",
RemotingHelper.exceptionSimpleDesc(e),
@@ -464,12 +450,10 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
returnType = ConsumeReturnType.TIME_OUT;
} else if (ConsumeConcurrentlyStatus.RECONSUME_LATER == status) {
returnType = ConsumeReturnType.FAILED;
- } else if (ConsumeConcurrentlyStatus.CONSUME_SUCCESS == status) {
- returnType = ConsumeReturnType.SUCCESS;
}
if
(ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.hasHook()) {
-
consumeMessageContext.getProps().put(MixAll.CONSUME_CONTEXT_TYPE,
returnType.name());
+
Objects.requireNonNull(consumeMessageContext).getProps().put(MixAll.CONSUME_CONTEXT_TYPE,
returnType.name());
}
if (null == status) {
@@ -481,7 +465,7 @@ public class ConsumeMessageConcurrentlyService implements
ConsumeMessageService
}
if
(ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.hasHook()) {
- consumeMessageContext.setStatus(status.toString());
+
Objects.requireNonNull(consumeMessageContext).setStatus(status.toString());
consumeMessageContext.setSuccess(ConsumeConcurrentlyStatus.CONSUME_SUCCESS ==
status);
ConsumeMessageConcurrentlyService.this.defaultMQPushConsumerImpl.executeHookAfter(consumeMessageContext);
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]