mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
[ROCKETMQ-30] Fixed method signature for Message Filter example and class loading from resources, closes apache/incubator-rocketmq#27
This commit is contained in:
@@ -16,6 +16,7 @@
|
||||
*/
|
||||
package org.apache.rocketmq.example.filter;
|
||||
|
||||
import java.io.File;
|
||||
import java.util.List;
|
||||
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
|
||||
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
|
||||
@@ -30,8 +31,11 @@ public class Consumer {
|
||||
public static void main(String[] args) throws InterruptedException, MQClientException {
|
||||
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("ConsumerGroupNamecc4");
|
||||
|
||||
String filterCode = MixAll.file2String("/home/admin/MessageFilterImpl.java");
|
||||
consumer.subscribe("TopicFilter7", "org.apache.rocketmq.example.filter.MessageFilterImpl",
|
||||
ClassLoader classLoader = Thread.currentThread().getContextClassLoader();
|
||||
File classFile = new File(classLoader.getResource("MessageFilterImpl.java").getFile());
|
||||
|
||||
String filterCode = MixAll.file2String(classFile);
|
||||
consumer.subscribe("TopicTest", "org.apache.rocketmq.example.filter.MessageFilterImpl",
|
||||
filterCode);
|
||||
|
||||
consumer.registerMessageListener(new MessageListenerConcurrently() {
|
||||
|
||||
@@ -17,13 +17,14 @@
|
||||
|
||||
package org.apache.rocketmq.example.filter;
|
||||
|
||||
import org.apache.rocketmq.common.filter.FilterContext;
|
||||
import org.apache.rocketmq.common.filter.MessageFilter;
|
||||
import org.apache.rocketmq.common.message.MessageExt;
|
||||
|
||||
public class MessageFilterImpl implements MessageFilter {
|
||||
|
||||
@Override
|
||||
public boolean match(MessageExt msg) {
|
||||
public boolean match(MessageExt msg, FilterContext context) {
|
||||
String property = msg.getProperty("SequenceId");
|
||||
if (property != null) {
|
||||
int id = Integer.parseInt(property);
|
||||
|
||||
Reference in New Issue
Block a user