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]

Reply via email to