2021-01-13 19:27发布
rocketMQ如何设置消费线程数
// 一个应用创建一个Consumer,由应用来维护此对象,可以设置为全局对象或者单例,ConsumerGroupName需要由应用来保证唯一
final DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(rocketMQConfig.getConsumerGroup());
consumer.setNamesrvAddr(rocketMQConfig.getNamesrvAddr());
// 订阅指定topic下tags分别等于TagA或TagC或TagD
//consumer.subscribe("Topic1", "TagA || TagC || TagD");
// 订阅指定topic下所有消息,注意:一个consumer对象可以订阅多个topic
consumer.subscribe(rocketMQConfig.getSmsReceiptTopic(), "*");
consumer.setVipChannelEnabled(false);
// 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费 如果非第一次启动,那么按照上次消费的位置继续消费
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
//单线程处理设置最大线程,最小线程也要设置,不然启动报错
consumer.setConsumeThreadMax(1);
consumer.setConsumeThreadMin(1);
//默认msgs里只有一条消息,可以通过设置consumeMessageBatchMaxSize参数来批量接收消息
//MessageListenerOrderly一个线程一个队列顺序接收
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
System.out.println(Thread.currentThread().getName() + " Receive New Messages: " + msgs);
MessageExt msg = (MessageExt) msgs.get(0);
String msgBody = new String(msg.getBody());
//短信回执主题
if (msg.getTopic().equals(rocketMQConfig.getSmsReceiptTopic())) {
log.info("消费者接收到短信回执消息:[{}]", msgBody);
httpSmsSvc.handleSmsReplyAndReport(msgBody);
} else {
log.error("出现意料之外的消息:[{}]", msg.toString());
}
return ConsumeOrderlyStatus.SUCCESS;
});
// Consumer对象在使用之前必须要调用start初始化,初始化一次即可
consumer.start();
System.out.println("Consumer Started.");
Runtime.getRuntime().addShutdownHook(new Thread(consumer::shutdown));
RocketMQ 版所提供的 TCP Java SDK 支持多线程消费,且适用于所有消息类型,本文介绍如何设置消费线程数的方法。
在启动 Consumer 时,设置一个 ConsumeThreadNums 属性即可。具体示例如下所示。
public static void main(String[] args) { Properties properties = new Properties(); properties.put(PropertyKeyConst.GROUP_ID, "GID_001"); properties.put(PropertyKeyConst.AccessKey, "xxxxxxxxxxxx"); properties.put(PropertyKeyConst.SecretKey, "xxxxxxxxxxxx"); /** * 设置消费端线程数固定为 20 */ properties.setProperty(PropertyKeyConst.ConsumeThreadNums,"20"); Consumer consumer =ONSFactory.createConsumer(properties); consumer.subscribe("TestTopic", "*", new MessageListener() { public Action consume(Message message, ConsumeContext context) { System.out.println("Receive: " + message); return Action.CommitMessage; } }); consumer.start(); System.out.println("Consumer Started"); }
Rocketmq消费分为push和pull两种方式,push为被动消费类型,pull为主动消费类型,push方式最终还是会从broker中pull消息。不同于pull的是,push首先要注册消费监听器,当监听器处触发后才开始消费消息,所以被称为“被动”消费。
public class PushConsumer {
public
class
PushConsumer {
static
void
main(String[] args)
throws
InterruptedException, MQClientException {
DefaultMQPushConsumer consumer =
new
DefaultMQPushConsumer(
"CID_JODIE_1"
);
consumer.subscribe(
"Jodie_topic_1023"
,
"*"
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
//wrong time format 2017_0422_221800
consumer.setConsumeTimestamp(
"20170422221800"
consumer.registerMessageListener(
MessageListenerConcurrently() {
@Override
ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeConcurrentlyContext context) {
System.out.printf(Thread.currentThread().getName() +
" Receive New Messages: "
+ msgs +
"%n"
return
ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
System.out.printf(
"Consumer Started.%n"
checkConfig 检查内容:
1.消费组 -- (不能与默认DEFAULT_CONSUMER同名)
2.消费模型 -- (默认CLUSTERING)
3.从何处开始消费 -- (默认CONSUME_FROM_LAST_OFFSET)
4.消费时间戳 -- (消息回溯,默认Default backtracking consumption time Half an hour ago)
5.消费负载均衡策略 -- (默认AllocateMessageQueueAveragely)
6.订阅关系 --(map类型,即可订阅多个topic;key=Topic, value=订阅描述)
7.消费监听 --(必须为orderly or concurrently类型之一)
8.消费消息的线程数量控制 -- (消费线程池最大、最小数量)
9.检查单队列并行消费允许的最大跨度 --(consumeConcurrentlyMaxSpan)
10.检查拉消息本地队列缓存消息最大数 --(pullThresholdForQueue)(processQueue.getMsgCount()记数)
11.检查拉取时间间隔 --(拉消息间隔,由于是长轮询,所以默认为0)
12.检查批量消费的个数 --(一次消费多少条消息)
13.检查批量拉取消息的个数 --(一次最多拉多少条)
Statement的execute(String query)方法用来执行任意的SQL查询,如果查询的结果是一个ResultSet,这个方法就返回true。如果结果不是ResultSet,比如insert或者update查询,它就会返回false。我们可以通过它的getResultSet方法来获取ResultSet,或者通过getUpda...
忙的时候项目期肯定要加班 但是每天加班应该还不至于
虽然Java人才越来越多,但是人才缺口也是很大的,我国对JAVA工程师的需求是所有软件工程师当中需求大的,达到全部需求量的60%-70%,所以Java市场在短时间内不可能饱和。其次,Java市场不断变化,人才需求也会不断增加。马云说过,未来的制造业要的不是石油,...
工信部证书含金量较高。工信部是国务院的下属结构,具有发放资质、证书的资格。其所发放的证书具有较强的权威性,在全国范围内收到认可,含金量通常都比较高。 工信部证书,其含义也就是工信部颁发并承认的某项技能证书,是具有法律效力的,并且是国家认可的...
学Java好不好找工作?看学完Java后能做些什么吧。一、大数据技术Hadoop以及其他大数据处理技术都是用Java或者其他,例如Apache的基于Java 的 HBase和Accumulo以及ElasticSearchas。但是Java在此领域并未占太大空间,但只要Hadoop和ElasticSearchas能够成长壮...
就是java的基础知识啊,比如Java 集合框架;Java 多线程;线程的五种状态;Java 虚拟机;MySQL (InnoDB);Spring 相关;计算机网络;MQ 消息队列诸如此类
#{}和${}这两个语法是为了动态传递参数而存在的,是Mybatis实现动态SQL的基础,总体上他们的作用是一致的(为了动态传参),但是在编译过程、是否自动加单引号、安全性、使用场景等方面有很多不同,下面详细比较两者间的区别:1.#{} 是 占位符 :动态解析 ...
没问题的,专科学历也能学习Java开发的,主要看自己感不感兴趣,只要认真学,市面上的培训机构不少都是零基础课程,能跟得上,或是自己先找些资料学习一下。
1、反射对单例模式的破坏采用反射的方式另辟蹊径实例了该类,导致程序中会存在不止一个实例。解决方案其思想就是采用一个全局变量,来标记是否已经实例化过了,如果已经实例化过了,第 二次实例化的时候,抛出异常2、clone()对单例模式的破坏当需要实现单例的...
优点: 一、实例控制 单例模式会阻止其他对象实例化其自己的单例对象的副本,从而确保所有对象都访问唯一实例。 二、灵活性 因为类控制了实例化过程,所以类可以灵活更改实例化过程。 缺点: 一、开销 虽然数量很少,但如果每次对象请求引用时都要...
这个主要是看你数组的长度是多少, 比如之前写过的一个程序有个数组存的是各个客户端的ip地址:string clientIp[4]={XXX, xxx, xxx, xxx};这个时候如果想把hash值对应到上面四个地址的话,就应该对4取余,这个时候p就应该为4...
哈希表的大小 · 关键字的分布情况 · 记录的查找频率 1.直接寻址法:取关键字或关键字的某个线性函数值为散列地址。即H(key)=key或H(key) = a·key + b,其中a和b为常数(这种散列函数叫做自身函数)。...
哈希表的大小取决于一组质数,原因是在hash函数中,你要用这些质数来做模运算(%)。而分析发现,如果不是用质数来做模运算的话,很多生活中的数据分布,会集中在某些点上。所以这里最后采用了质数做模的除数。 因为用质数做了模的除数,自然存储空间的大小也用质数了...
是啊,哈希函数的设计至关重要,好的哈希函数会尽可能地保证计算简单和散列地址分布均匀,但是,我们需要清楚的是,数组是一块连续的固定长度的内存空间
解码查表优化算法,seo优化
1.对对象元素中的关键字(对象中的特有数据),进行哈希算法的运算,并得出一个具体的算法值,这个值 称为哈希值。2.哈希值就是这个元素的位置。3.如果哈希值出现冲突,再次判断这个关键字对应的对象是否相同。如果对象相同,就不存储,因为元素重复。如果对象不同,就...
最多设置5个标签!
// 一个应用创建一个Consumer,由应用来维护此对象,可以设置为全局对象或者单例,ConsumerGroupName需要由应用来保证唯一
final DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(rocketMQConfig.getConsumerGroup());
consumer.setNamesrvAddr(rocketMQConfig.getNamesrvAddr());
// 订阅指定topic下tags分别等于TagA或TagC或TagD
//consumer.subscribe("Topic1", "TagA || TagC || TagD");
// 订阅指定topic下所有消息,注意:一个consumer对象可以订阅多个topic
consumer.subscribe(rocketMQConfig.getSmsReceiptTopic(), "*");
consumer.setVipChannelEnabled(false);
// 设置Consumer第一次启动是从队列头部开始消费还是队列尾部开始消费 如果非第一次启动,那么按照上次消费的位置继续消费
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET);
//单线程处理设置最大线程,最小线程也要设置,不然启动报错
consumer.setConsumeThreadMax(1);
consumer.setConsumeThreadMin(1);
//默认msgs里只有一条消息,可以通过设置consumeMessageBatchMaxSize参数来批量接收消息
//MessageListenerOrderly一个线程一个队列顺序接收
consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) -> {
System.out.println(Thread.currentThread().getName() + " Receive New Messages: " + msgs);
MessageExt msg = (MessageExt) msgs.get(0);
String msgBody = new String(msg.getBody());
//短信回执主题
if (msg.getTopic().equals(rocketMQConfig.getSmsReceiptTopic())) {
log.info("消费者接收到短信回执消息:[{}]", msgBody);
httpSmsSvc.handleSmsReplyAndReport(msgBody);
} else {
log.error("出现意料之外的消息:[{}]", msg.toString());
}
return ConsumeOrderlyStatus.SUCCESS;
});
// Consumer对象在使用之前必须要调用start初始化,初始化一次即可
consumer.start();
System.out.println("Consumer Started.");
Runtime.getRuntime().addShutdownHook(new Thread(consumer::shutdown));
RocketMQ 版所提供的 TCP Java SDK 支持多线程消费,且适用于所有消息类型,本文介绍如何设置消费线程数的方法。
在启动 Consumer 时,设置一个 ConsumeThreadNums 属性即可。具体示例如下所示。
Rocketmq消费分为push和pull两种方式,push为被动消费类型,pull为主动消费类型,push方式最终还是会从broker中pull消息。不同于pull的是,push首先要注册消费监听器,当监听器处触发后才开始消费消息,所以被称为“被动”消费。
public
class
PushConsumer {
public
static
void
main(String[] args)
throws
InterruptedException, MQClientException {
DefaultMQPushConsumer consumer =
new
DefaultMQPushConsumer(
"CID_JODIE_1"
);
consumer.subscribe(
"Jodie_topic_1023"
,
"*"
);
consumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
//wrong time format 2017_0422_221800
consumer.setConsumeTimestamp(
"20170422221800"
);
consumer.registerMessageListener(
new
MessageListenerConcurrently() {
@Override
public
ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeConcurrentlyContext context) {
System.out.printf(Thread.currentThread().getName() +
" Receive New Messages: "
+ msgs +
"%n"
);
return
ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
System.out.printf(
"Consumer Started.%n"
);
}
}
checkConfig 检查内容:
1.消费组 -- (不能与默认DEFAULT_CONSUMER同名)
2.消费模型 -- (默认CLUSTERING)
3.从何处开始消费 -- (默认CONSUME_FROM_LAST_OFFSET)
4.消费时间戳 -- (消息回溯,默认Default backtracking consumption time Half an hour ago)
5.消费负载均衡策略 -- (默认AllocateMessageQueueAveragely)
6.订阅关系 --(map类型,即可订阅多个topic;key=Topic, value=订阅描述)
7.消费监听 --(必须为orderly or concurrently类型之一)
8.消费消息的线程数量控制 -- (消费线程池最大、最小数量)
9.检查单队列并行消费允许的最大跨度 --(consumeConcurrentlyMaxSpan)
10.检查拉消息本地队列缓存消息最大数 --(pullThresholdForQueue)(processQueue.getMsgCount()记数)
11.检查拉取时间间隔 --(拉消息间隔,由于是长轮询,所以默认为0)
12.检查批量消费的个数 --(一次消费多少条消息)
13.检查批量拉取消息的个数 --(一次最多拉多少条)
相关问题推荐
Statement的execute(String query)方法用来执行任意的SQL查询,如果查询的结果是一个ResultSet,这个方法就返回true。如果结果不是ResultSet,比如insert或者update查询,它就会返回false。我们可以通过它的getResultSet方法来获取ResultSet,或者通过getUpda...
忙的时候项目期肯定要加班 但是每天加班应该还不至于
虽然Java人才越来越多,但是人才缺口也是很大的,我国对JAVA工程师的需求是所有软件工程师当中需求大的,达到全部需求量的60%-70%,所以Java市场在短时间内不可能饱和。其次,Java市场不断变化,人才需求也会不断增加。马云说过,未来的制造业要的不是石油,...
工信部证书含金量较高。工信部是国务院的下属结构,具有发放资质、证书的资格。其所发放的证书具有较强的权威性,在全国范围内收到认可,含金量通常都比较高。 工信部证书,其含义也就是工信部颁发并承认的某项技能证书,是具有法律效力的,并且是国家认可的...
学Java好不好找工作?看学完Java后能做些什么吧。一、大数据技术Hadoop以及其他大数据处理技术都是用Java或者其他,例如Apache的基于Java 的 HBase和Accumulo以及ElasticSearchas。但是Java在此领域并未占太大空间,但只要Hadoop和ElasticSearchas能够成长壮...
就是java的基础知识啊,比如Java 集合框架;Java 多线程;线程的五种状态;Java 虚拟机;MySQL (InnoDB);Spring 相关;计算机网络;MQ 消息队列诸如此类
#{}和${}这两个语法是为了动态传递参数而存在的,是Mybatis实现动态SQL的基础,总体上他们的作用是一致的(为了动态传参),但是在编译过程、是否自动加单引号、安全性、使用场景等方面有很多不同,下面详细比较两者间的区别:1.#{} 是 占位符 :动态解析 ...
没问题的,专科学历也能学习Java开发的,主要看自己感不感兴趣,只要认真学,市面上的培训机构不少都是零基础课程,能跟得上,或是自己先找些资料学习一下。
1、反射对单例模式的破坏采用反射的方式另辟蹊径实例了该类,导致程序中会存在不止一个实例。解决方案其思想就是采用一个全局变量,来标记是否已经实例化过了,如果已经实例化过了,第 二次实例化的时候,抛出异常2、clone()对单例模式的破坏当需要实现单例的...
优点: 一、实例控制 单例模式会阻止其他对象实例化其自己的单例对象的副本,从而确保所有对象都访问唯一实例。 二、灵活性 因为类控制了实例化过程,所以类可以灵活更改实例化过程。 缺点: 一、开销 虽然数量很少,但如果每次对象请求引用时都要...
这个主要是看你数组的长度是多少, 比如之前写过的一个程序有个数组存的是各个客户端的ip地址:string clientIp[4]={XXX, xxx, xxx, xxx};这个时候如果想把hash值对应到上面四个地址的话,就应该对4取余,这个时候p就应该为4...
哈希表的大小 · 关键字的分布情况 · 记录的查找频率 1.直接寻址法:取关键字或关键字的某个线性函数值为散列地址。即H(key)=key或H(key) = a·key + b,其中a和b为常数(这种散列函数叫做自身函数)。...
哈希表的大小取决于一组质数,原因是在hash函数中,你要用这些质数来做模运算(%)。而分析发现,如果不是用质数来做模运算的话,很多生活中的数据分布,会集中在某些点上。所以这里最后采用了质数做模的除数。 因为用质数做了模的除数,自然存储空间的大小也用质数了...
是啊,哈希函数的设计至关重要,好的哈希函数会尽可能地保证计算简单和散列地址分布均匀,但是,我们需要清楚的是,数组是一块连续的固定长度的内存空间
解码查表优化算法,seo优化
1.对对象元素中的关键字(对象中的特有数据),进行哈希算法的运算,并得出一个具体的算法值,这个值 称为哈希值。2.哈希值就是这个元素的位置。3.如果哈希值出现冲突,再次判断这个关键字对应的对象是否相同。如果对象相同,就不存储,因为元素重复。如果对象不同,就...