Skip to content

Commit e1ab73f

Browse files
joeweejoewee
andauthored
[ISSUE #429]Use 'consumeThreadNumber' instead of 'consumeThreadMax' (#431)
Co-authored-by: joewee <371499220@qq.com>
1 parent 45833f6 commit e1ab73f

3 files changed

Lines changed: 25 additions & 7 deletions

File tree

rocketmq-spring-boot/src/main/java/org/apache/rocketmq/spring/annotation/RocketMQMessageListener.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import java.lang.annotation.Retention;
2323
import java.lang.annotation.RetentionPolicy;
2424
import java.lang.annotation.Target;
25+
import java.util.concurrent.LinkedBlockingQueue;
2526

2627
@Target(ElementType.TYPE)
2728
@Retention(RetentionPolicy.RUNTIME)
@@ -72,9 +73,19 @@
7273

7374
/**
7475
* Max consumer thread number.
76+
* @deprecated This property is not work well, because the consumer thread pool executor use
77+
* {@link LinkedBlockingQueue} with default capacity bound (Integer.MAX_VALUE), use
78+
* {@link RocketMQMessageListener#consumeThreadNumber} .
79+
* @see <a href="https://github.com/apache/rocketmq-spring/issues/429">issues#429</a>
7580
*/
81+
@Deprecated
7682
int consumeThreadMax() default 64;
7783

84+
/**
85+
* consumer thread number.
86+
*/
87+
int consumeThreadNumber() default 20;
88+
7889
/**
7990
* Max re-consume times.
8091
*

rocketmq-spring-boot/src/main/java/org/apache/rocketmq/spring/support/DefaultRocketMQListenerContainer.java

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,8 @@ public class DefaultRocketMQListenerContainer implements InitializingBean,
105105

106106
private int consumeThreadMax = 64;
107107

108+
private int consumeThreadNumber = 20;
109+
108110
private String charset = "UTF-8";
109111

110112
private MessageConverter messageConverter;
@@ -186,6 +188,10 @@ public int getConsumeThreadMax() {
186188
return consumeThreadMax;
187189
}
188190

191+
public int getConsumeThreadNumber() {
192+
return consumeThreadNumber;
193+
}
194+
189195
public String getCharset() {
190196
return charset;
191197
}
@@ -227,7 +233,8 @@ public void setRocketMQMessageListener(RocketMQMessageListener anno) {
227233
this.rocketMQMessageListener = anno;
228234

229235
this.consumeMode = anno.consumeMode();
230-
this.consumeThreadMax = anno.consumeThreadMax();
236+
this.consumeThreadMax = anno.consumeThreadNumber();
237+
this.consumeThreadNumber = anno.consumeThreadNumber();
231238
this.messageModel = anno.messageModel();
232239
this.selectorType = anno.selectorType();
233240
this.selectorExpression = anno.selectorExpression();
@@ -612,10 +619,9 @@ private void initRocketMQPushConsumer() throws MQClientException {
612619
if (accessChannel != null) {
613620
consumer.setAccessChannel(accessChannel);
614621
}
615-
consumer.setConsumeThreadMax(consumeThreadMax);
616-
if (consumeThreadMax < consumer.getConsumeThreadMin()) {
617-
consumer.setConsumeThreadMin(consumeThreadMax);
618-
}
622+
//set the consumer core thread number and maximum thread number has the same value
623+
consumer.setConsumeThreadMax(consumeThreadNumber);
624+
consumer.setConsumeThreadMin(consumeThreadNumber);
619625
consumer.setConsumeTimeout(consumeTimeout);
620626
consumer.setMaxReconsumeTimes(maxReconsumeTimes);
621627
switch (messageModel) {

rocketmq-spring-boot/src/test/java/org/apache/rocketmq/spring/support/DefaultRocketMQListenerContainerTest.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -244,7 +244,8 @@ public void testSetRocketMQMessageListener() {
244244
container.setRocketMQMessageListener(anno);
245245

246246
assertEquals(anno.consumeMode(), container.getConsumeMode());
247-
assertEquals(anno.consumeThreadMax(), container.getConsumeThreadMax());
247+
assertEquals(anno.consumeThreadNumber(), container.getConsumeThreadMax());
248+
assertEquals(anno.consumeThreadNumber(), container.getConsumeThreadNumber());
248249
assertEquals(anno.messageModel(), container.getMessageModel());
249250
assertEquals(anno.selectorType(), container.getSelectorType());
250251
assertEquals(anno.selectorExpression(), container.getSelectorExpression());
@@ -256,7 +257,7 @@ public void testSetRocketMQMessageListener() {
256257

257258
@RocketMQMessageListener(consumerGroup = "abc1", topic = "test",
258259
consumeMode = ConsumeMode.ORDERLY,
259-
consumeThreadMax = 3456,
260+
consumeThreadNumber = 3456,
260261
messageModel = MessageModel.BROADCASTING,
261262
selectorType = SelectorType.SQL92,
262263
selectorExpression = "selectorExpression",

0 commit comments

Comments
 (0)