Skip to content

Commit b62097b

Browse files
committed
support simple consumer supports subscribing to multiple topics
1 parent 4075513 commit b62097b

5 files changed

Lines changed: 38 additions & 36 deletions

File tree

rocketmq-v5-client-spring-boot-samples/rocketmq-v5-client-consume-simple-subscribe-muliti-topic-demo/src/main/resources/application.properties

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -14,10 +14,11 @@
1414
# limitations under the License.
1515

1616
rocketmq.simple-consumer.endpoints=localhost:8081
17-
rocketmq.simple-consumer.consumer-group=demo-group
18-
rocketmq.simple-consumer.filter-expression-map.demo-topic1.tag=tagA
19-
rocketmq.simple-consumer.filter-expression-map.demo-topic1.filter-expression-type=tag
20-
rocketmq.simple-consumer.filter-expression-map.demo-topic2.tag=tagB
21-
rocketmq.simple-consumer.filter-expression-map.demo-topic2.filter-expression-type=tag
17+
rocketmq.simple-consumer.consumer-group=localhost:8081
18+
rocketmq.simple-consumer.subscription-expressions.demo-topic.tag=tagA
19+
rocketmq.simple-consumer.subscription-expressions.demo-topic.filter-expression-type=tag
20+
rocketmq.simple-consumer.subscription-expressions.demo-topic2.tag=tagB
21+
rocketmq.simple-consumer.subscription-expressions.demo-topic2.filter-expression-type=tag
2222
#rocketmq.simple-consumer.access-key=
2323
#rocketmq.simple-consumer.secret-key=
24+

rocketmq-v5-client-spring-boot/src/main/java/org/apache/rocketmq/client/autoconfigure/RocketMQAutoConfiguration.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -118,13 +118,13 @@ public SimpleConsumerBuilder simpleConsumerBuilder(RocketMQProperties rocketMQPr
118118
}
119119

120120
// Set the subscription for the consumer.
121-
if (simpleConsumer.getFilterExpressionMap().isEmpty()) {
121+
if (simpleConsumer.getSubscriptionExpressions().isEmpty()) {
122122
FilterExpression filterExpression = RocketMQUtil.createFilterExpression(simpleConsumer.getTag(), simpleConsumer.getFilterExpressionType());
123123
if (Objects.nonNull(filterExpression)) {
124124
simpleConsumerBuilder.setSubscriptionExpressions(Collections.singletonMap(simpleConsumer.getTopic(), filterExpression));
125125
}
126126
} else {
127-
Map<String, FilterExpression> subscriptionExpressions = RocketMQUtil.createSubscriptionExpressions(simpleConsumer.getFilterExpressionMap());
127+
Map<String, FilterExpression> subscriptionExpressions = RocketMQUtil.createSubscriptionExpressions(simpleConsumer.getSubscriptionExpressions());
128128
simpleConsumerBuilder.setSubscriptionExpressions(subscriptionExpressions);
129129
}
130130
return simpleConsumerBuilder;

rocketmq-v5-client-spring-boot/src/main/java/org/apache/rocketmq/client/autoconfigure/RocketMQProperties.java

Lines changed: 6 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
*/
1717
package org.apache.rocketmq.client.autoconfigure;
1818

19+
import org.apache.rocketmq.client.common.FilterExpression;
1920
import org.springframework.boot.context.properties.ConfigurationProperties;
2021

2122
import java.util.Map;
@@ -216,7 +217,7 @@ public static class SimpleConsumer {
216217
/**
217218
* key is topic
218219
*/
219-
private Map<String, FilterExpression> filterExpressionMap;
220+
private Map<String, FilterExpression> subscriptionExpressions;
220221

221222
public String getAccessKey() {
222223
return accessKey;
@@ -306,12 +307,12 @@ public void setNamespace(String namespace) {
306307
this.namespace = namespace;
307308
}
308309

309-
public Map<String, FilterExpression> getFilterExpressionMap() {
310-
return filterExpressionMap;
310+
public Map<String, FilterExpression> getSubscriptionExpressions() {
311+
return subscriptionExpressions;
311312
}
312313

313-
public void setFilterExpressionMap(Map<String, FilterExpression> filterExpressionMap) {
314-
this.filterExpressionMap = filterExpressionMap;
314+
public void setSubscriptionExpressions(Map<String, FilterExpression> subscriptionExpressions) {
315+
this.subscriptionExpressions = subscriptionExpressions;
315316
}
316317

317318
@Override
@@ -329,27 +330,4 @@ public String toString() {
329330
'}';
330331
}
331332
}
332-
333-
public static class FilterExpression {
334-
private String tag;
335-
336-
private String filterExpressionType;
337-
338-
public String getTag() {
339-
return tag;
340-
}
341-
342-
public void setTag(String tag) {
343-
this.tag = tag;
344-
}
345-
346-
public String getFilterExpressionType() {
347-
return filterExpressionType;
348-
}
349-
350-
public void setFilterExpressionType(String filterExpressionType) {
351-
this.filterExpressionType = filterExpressionType;
352-
}
353-
}
354-
355333
}
Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,23 @@
1+
package org.apache.rocketmq.client.common;
2+
3+
public class FilterExpression {
4+
private String tag;
5+
6+
private String filterExpressionType;
7+
8+
public String getTag() {
9+
return tag;
10+
}
11+
12+
public void setTag(String tag) {
13+
this.tag = tag;
14+
}
15+
16+
public String getFilterExpressionType() {
17+
return filterExpressionType;
18+
}
19+
20+
public void setFilterExpressionType(String filterExpressionType) {
21+
this.filterExpressionType = filterExpressionType;
22+
}
23+
}

rocketmq-v5-client-spring-boot/src/main/java/org/apache/rocketmq/client/support/RocketMQUtil.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -187,7 +187,7 @@ public static FilterExpression createFilterExpression(String tag, String type) {
187187
return filterExpression;
188188
}
189189

190-
public static Map<String, FilterExpression> createSubscriptionExpressions(Map<String, RocketMQProperties.FilterExpression> map) {
190+
public static Map<String, FilterExpression> createSubscriptionExpressions(Map<String, org.apache.rocketmq.client.common.FilterExpression> map) {
191191
Map<String, FilterExpression> subscriptionExpressions = new HashMap<>();
192192
map.forEach((topic, expression) -> {
193193
FilterExpressionType filterExpressionType = "tag".equalsIgnoreCase(expression.getFilterExpressionType()) ? FilterExpressionType.TAG : FilterExpressionType.SQL92;

0 commit comments

Comments
 (0)