性能优化-请求优先级
业务系统中,针对多种业务请求,期望高优先级的任务尽响应客户端,低优先级的任务在没有高优先级任务情况下可以充分得到执行,有高优先级任务情况下偶尔执行或者不执行。这种场景有
- 物联网系统:设备上下线的任务处理优先级高,其次是设备故障告警事件优先级高,而设备属性上报优先级低
- 电商系统,支付业务是核心业务,需要高优先级执行;电商中请求标记为VIP的用户,需要优先得到系统的服务
- 金融系统中,风控系统是核心业务,需要高优先级执行;金融系统功能中,低频的大额资金业务比高频的小金额业务更优先执行。
实现请求的优先级处理,通常有如下步骤
- 标记任务为优先级,如物联网设备上线业务的优先级高于属性上报优先级
- 为不同优先级的请求分配不同的资源处理,如果是分布式系统,可以通过代理(api网关),或者微服务框架,消息中间件等,派发不同优先级任务到不同的服务集群。如果是单体系统,高优先级任务应该得到线程优先执行。
下图是在分布式系统中,采用消息中间件Kafka派发不同优先级任务到不同的Topic以实现优先级处理

Producer根据业务类型,派发不同功能优先级的业务到Topic里,如上图建立了俩个Topic,一个是高优先级,一个是低优先级。每个Topic都配置不同的Consumer集群 进行处理。比如大额交易放入高优先级队列,此时尽管对应的Consumer部署较少节点,但因为大额交易请求数量较少,能得到优先执行。 小额交易进入低优先级Topic,然后派发到低优先级Consumer集群处理。
Kafka使用了多个Topic实现优先级,RabbitMQ,Apache ActiveMQ以及商业消息中间件支持 优先级队列,使用一个队列,其能对消息进行排序,供消费者按照优先级处理。
单体系统中,也可以通过使用多个线程池实现优先级处理。比如设定设备上下线有一个线程池,其他业务用另外一个线程池
ThreadPoolExecutor onlineEventPool = ...
ThreadPoolExecutor warnningEvnentPool = ...
ThreadPoolExecutor nomalEventPool = ...分布式下使用多个集群,或者单体系统使用多个线程池,除了简单快速实现优先级处理外,还有个优势是能“物理隔离”优先级任务到不同环境执行。比如上图Kafka实现优先级任务处理,单独维护低优先级Consumer,不影响高优先级Consumer的执行。
但这种方式在特定场景下有俩个缺点需要解决:
- 为不同优先级创建不同的集群资源或者线程池,但如果设定优先级较多,需要每一个优先级创建一个集群资源或者线程池。
- 即使较少优先级设定对应了较少的集群资源或者线程池,也可能造成资源浪费,比如专门处理高优先级的集群资源长期处于不繁忙状态,不利于降本增效。
解决方案是只创建一个集群或者一个线程池,通过一个“派发器” ,将高优先级任务优先转发到唯一的一个集群或者线程池,低优先级先暂时存放到一个队列等待派发。
单体系统最简单方法是使用PriorityBlockingQueue 构造ThreadPoolExecutor,PriorityBlockingQueue 在构造时候要求实现Comparator以计算优先级,Comparator::compare 返回数值大于0将排在队列前优先执行,小于0将延后执行
public PriorityBlockingQueue(int initialCapacity,
Comparator<? super E> comparator)需要注意的是,如果高优先级任务总是得到优先处理,会导致低优先级任务没有机会执行(需要结合实际业务考虑是否会发生这种情况,比如大多数高优先级任务个数是低于普通优先级的,不太可能存在此情况),因此队列的优先级通常不是直接根据业务优先级计算而得,而是综合了业务优先级和在队列等待时间计算,如果低优先级的任务在队列等待时间较长,其优先级提升,比如DelegateTask的getPriority,用于动态计算在队列中的优先级,当等待时间超过500毫秒,优先级得到提高。
public class DelegateTask implements Runnable {
private final int priority;
private final Runnable delegate;
public final long enqueueTs;
/** 老化系数:每等待多少毫秒,+1老化分数;数值越大老化越慢 */
private static final long AGING_UNIT_MS = 500;
public DelegateTask(int priority, Runnable delegate) {
this.priority = priority;
this.delegate = delegate;
this.enqueueTs = System.currentTimeMillis();
}
/**
* 计算队列中的动态优先级,PriorityBlockingQueue的Comparator类,将使用getPriority 用于计算优先级
*/
public int getPriority() {
long waitMs = System.currentTimeMillis() - enqueueTs;
int agingScore = (int) (waitMs / AGING_UNIT_MS);
return priority + agingScore;
}
@Override
public void run() {
delegate.run();
}
}另外一个方法是实现MultipleQueue,作为参数传入线程池,MultipleQueue内部包含了多个BlockingQueue,代表不同的优先级。当ThreadPoolExecutor 调用 MultipleQueue::poll 取出一个等待执行任务时候,offer方法需调用Policy类来决定返回一个任务
public class MultipleQueue implements Queue {
BlockingQueue highQueue;
BlockingQueue normalQueue;
Policy policy;
public MultipleQueue(BlockingQueue highQueue,BlockingQueue normalQueue,Policy policy){
}
//线程次内部将调用poll方法获取任务,这里重新实现,通过policy决定返回哪个队列的任务
@Override
public Object poll() {
return policy.get(highQueue,normalQueue);
}
}MultipleQueue 维护多个队列,每个队列代表了一个优先级。当ThreadPoolExecutor调用poll方法获取一个可执行的任务时候,Policy来决定返回哪个队列的的任务,如下类AlwayHightPolicy总是优先返回高优先级任务
public class AlwayHightPolicy implements Policy{
@Override
public Object get(BlockingQueue<?> high, BlockingQueue<?> nomal) {
//总是优先执行高优先级任务
if(!high.isEmpty()){return high.poll();}
return nomal.poll();
}
}可以引入计数器,实现权重,当从高优先级队列中获取4个任务时候,可以从低优先级队列获取一个任务,当没有高优先级任务时候,总是取出低优先级任务。
public class WeightPolicy implements Policy{
int highCount=0;
int weight =4;
@Override
public Object get(BlockingQueue<?> highQueue, BlockingQueue<?> normalQueue) {
if(highCount%weight!=(weight-1)){
getHight(highQueue,normalQueue);
}
return getNormal(highQueue,normalQueue);
}
//获取高优先级队列任务
protected Object getHight(BlockingQueue<?> highQueue, BlockingQueue<?> normalQueue){
if(!highQueue.isEmpty()){
highCount++;
return highQueue.poll();
}else{
return normalQueue.poll();
}
}
//获取低优先级任务
protected Object getNormal(BlockingQueue<?> highQueue, BlockingQueue<?> normalQueue){
if(!normalQueue.isEmpty()){
return normalQueue.poll();
}else{
highCount++;
return highQueue.poll();
}
}
}MultipleQueue 可以扩展支持多个队列,每个队列代表了不同优先级。参考本书附带的例子。
