解锁新姿势 | 如何用配置中心实现全局动态流控?

简介: 当资源成为瓶颈时,服务框架需要对消费者做限流,启动流控保护机制。流量控制有多种策略,比较常用的有:针对访问速率的静态流控、针对资源占用的动态流控、针对消费者并发连接数的连接控制和针对并行访问数的并发控制。在实践中,各种流量控制策略需要综合使用才能起到较好的效果。

在分布式架构中,应用和应用之间的调用类型分为以下两种,流控方式也略有不同。

同步RPC类调用,比如RESTful,Dubbo,HSF等都属于该类。对于该类同步调用,通常限流方式为两种:针对服务提供者的并发全局流控,或针对服务消费者的并发局部流控。两种的控制手段类似,都是通过限制服务端或客服端并发调用数来进行限制。

异步MQ类调用,典型如RocketMQ,      Kafka,等。对于该类异步调用,通常限流方式是在订阅端限流。限流方式为两种:针对消息订阅者的并发流控,或针对消息订阅者的消费延时流控。

针对消息订阅者的消费延时流控基本原理是,在每次客户端消费时,可以增加一个延时来控制消费速度,这样理论消费并发最快速度为:

MaxRate = 1 / ConsumInterval * ConcurrentThreadNumber

比如如果消息并发消费线程为20,延时为100ms,则理论上可以将并发消费控制在200以下。具体公式如下:

200 = 1 / 0.1 * 20

相比并发线程数流控,消费延时流控优点在于实现相对简单,对MQ类客户端包依赖较少,不需要客户端提供控制并发线程数的动态调整接口。

以上各种流量控制方法,在分布式架构下,如果要做到全局动态控制,一个简单的技术方法是依赖配置中心,即通过配置中心来进行流控参数的下发。

下面章节详细介绍如何基于配置中心来实现异步消息消费的全局动态流控。使用的例子为阿里云上的 MQ (消息队列)和 ACM (应用配置管理)两款产品。

注:之所以用MQ为示例是因为在本文撰写之时,正好MQ Consumer Client SDK并不支持动态调整现成并发数,因此通过基于ACM来动态调整消费延迟的方法正好可以解决MQ消费流控动态的问题。

基于消费延时流控的基本原理

基本原理如下。其中,管理员或应用程序通过ACM控制台发布消费延时配置(RCV_INTERVAL_TIME),所有MQ消费程序订阅该配置。理论上,该配置从发布到下发所有客户端,可以在1秒内完成(取决于网络演示)。

8a6af60491405721f48152e02dd48061281bbd73

代码示例

该章节基于配置中心来实现异步消息消费的全局动态流控的代码示例。使用的例子为阿里云上的MQ(消息队列)和ACM(应用配置管理)两款产品,基于Java语言。关于SDK的详细介绍,可参见两款产品的官方文档。

在ACM上创建消费延时的参数,截屏如下。

d77e20cf46705d493c022f52e87526394b56350b

设置全局消费延时变量

首先,设置消费接收延时的全局变量, 如下。

  // 初始化消息接收延时参数,单位为millisecond

            static int RCV_INTERVAL_TIME = 10000;

            // 初始化配置服务,控制台通过示例代码自动获取下面参数

            ConfigService.init("acm.aliyun.com", /*租户ID*/"xxx", /*AK*/"xxx", /*SK*/"yyy");    

            // 主动获取配置

            String content = ConfigService.getConfig("app.mq.qos", "DEFAULT_GROUP", 6000);

            Properties p = new Properties();

            try {

                p.load(new StringReader(content));

                RCV_INTERVAL_TIME = Integer.valueOf(p.getProperty("RCV_INTERVAL_TIME"));

            } catch (IOException e) {

                e.printStackTrace();

            }

其次,设置ACM listener,确保当配置被修改时,即使更新 RCV_INTERVAL_TIME 参数, 如下。

// 初始化的时候,给配置添加监听,配置变更会回调通知

            ConfigService.addListener("app.mq.qos", "DEFAULT_GROUP", new ConfigChangeListener() {

                public void receiveConfigInfo(String configInfo) {

                    Properties p = new Properties();

                    try {

                        p.load(new StringReader(configInfo));

                        RCV_INTERVAL_TIME = Integer.valueOf(p.getProperty("RCV_INTERVAL_TIME"));

                    } catch (IOException e) {

                        e.printStackTrace();

                    }

                }

            });

设置 MQ 消费延时逻辑


完整实例如下。

注:这里 RCV_INTERVAL_TIME 参数的访问是故意没有加锁的,读者可以自行思考原因。Aliyun ONS Client不提供动态线程并发数,默认并发为20。因此这里正好使用消费延时参数来动态调节QoS。

//以下代码可直接贴在Main()函数里

        Properties properties = new Properties();

        properties.put(PropertyKeyConst.ConsumerId, "CID_consumer_group");

        properties.put(PropertyKeyConst.AccessKey,"xxx");

        properties.put(PropertyKeyConst.SecretKey, "yyy");

        properties.setProperty(PropertyKeyConst.SendMsgTimeoutMillis, "3000");

        // 设置 TCP 接入域名(此处以公共云生产环境为例)

        properties.put(PropertyKeyConst.ONSAddr,

          "http://onsaddr-internet.aliyun.com/rocketmq/nsaddr4client-internet");

        Consumer consumer = ONSFactory.createConsumer(properties);

        consumer.subscribe(/*Topic*/"topic-name", /*Tag*/null, new MessageListener() 

        {

            public Action consume(Message message, ConsumeContext context) {

                // MQ Subscribe QoS logical start, 

                // Each consuming process will sleep for RCV_INTERVAL_TIME seconds with 100 ms sleeping cycle.

                // Within each cycle, the thread will check RCV_INTERVAL_TIME in case it's set to a smaller value. 

                // RCV_INTERVAL_TIME <= 0 means no sleeping.

                int rcvIntervalTimeLeft = RCV_INTERVAL_TIME;

                while (rcvIntervalTimeLeft > 0) {

                    if (rcvIntervalTimeLeft > RCV_INTERVAL_TIME) {

                        rcvIntervalTimeLeft = RCV_INTERVAL_TIME;

                    }

                    try {

                        if (rcvIntervalTimeLeft >= 100) {

                            rcvIntervalTimeLeft -= 100;

                            Thread.sleep(100);

                        } else {

                            Thread.sleep(rcvIntervalTimeLeft);

                            rcvIntervalTimeLeft = 0;

                        }

                    } catch (InterruptedException e) {

                        e.printStackTrace();

                    }

                }

                // MQ Subscribe interval logical ends

                System.out.println("Receive: " + message);

                /*

                 * Put your business logic here.

                 */

                doSomething();

                return Action.CommitMessage;

            }

        });

        consumer.start();

运行结果


单机运行consumer进行消费,假设queue内的消息无限多,不存在消费万的情况,分三段测试,分别运行约5分钟,通过ACM配置推送来达到以下效果。

RCV_INTERVAL_TIME      = 100 ms

RCV_INTERVAL_TIME      = 5000 ms

RCV_INTERVAL_TIME      = 1000 ms

结果如下,在单MQ消费业务处理耗时约100ms情况下的,单机并发20线程的测试结果。

RCV_INTERVAL_TIME  = 100 ms:平均消费性能约为 9000 tpm 左右

RCV_INTERVAL_TIME  = 5000 ms:平均消费性能被限制到了 200 tpm 左右

RCV_INTERVAL_TIME  = 1000 ms:平均消费性能回升到到了 1100 tpm 左右

以上结果基本达到消费和 tpm 成反比的预期,最关键的是整个过程中,应用不中断,流控推送结果秒级生效到分布式集群。单机性能结果如下所示。

31209e6536769067a959296f4a82122945984982

相关产品详情请参见:

  • 消息产品

Aliyun MQ:aliyun.com/product/ons

  • 配置中心产品

Aliyun ACM:aliyun.com/product/acm




原文发布时间为:2018-01-19

本文作者:杨奕

本文来自云栖社区合作伙伴“阿里技术”,了解相关信息可以关注“阿里技术”微信公众号

相关实践学习
快速体验阿里云云消息队列RocketMQ版
本实验将带您快速体验使用云消息队列RocketMQ版Serverless系列实例进行获取接入点、创建Topic、创建订阅组、收发消息、查看消息轨迹和仪表盘。
消息队列 MNS 入门课程
1、消息队列MNS简介 本节课介绍消息队列的MNS的基础概念 2、消息队列MNS特性 本节课介绍消息队列的MNS的主要特性 3、MNS的最佳实践及场景应用 本节课介绍消息队列的MNS的最佳实践及场景应用案例 4、手把手系列:消息队列MNS实操讲 本节课介绍消息队列的MNS的实际操作演示 5、动手实验:基于MNS,0基础轻松构建 Web Client 本节课带您一起基于MNS,0基础轻松构建 Web Client
相关文章
|
4天前
|
人工智能 API 内存技术
刚刚 DeepSeek V4.1 Flash 开启内测,1 分钟教你用上!
刚刚 DeepSeek 内测群发布了 DeepSeek V4.1 Flash 中间版本内测的消息,这次的模型采用了新的结构,原生支持多模态、能力更强、速度更快、且成本更低。
1659 5
|
8天前
|
人工智能 运维 BI
阿里云千问办公QwenWork深度解析:基于Qwen3.8,六大核心能力重构企业全自动化工作流与计费选型指南
传统AI办公工具大多停留在对话问答、文档摘要、简单文案生成层面,只能完成单点碎片化任务,无法自主拆解复杂业务流程,很难串联多工具、多文档、外部业务系统完成端到端完整工作交付。很多企业在落地AI办公的时候,需要组合多款不同工具,来回切换界面,手动复制粘贴中间结果,智能化改造落地门槛居高不下。千问办公QwenWork是整合多款智能体产品能力打造的一体化企业办公智能体平台,底层基座依托Qwen3.8大模型,打通桌面端Agent、云端Agent、企业协同Agent三种运行形态,不再局限简单问答,接收业务目标之后自主拆解任务步骤,调用各类工具,处理文档、表格、浏览器自动化、数据查询,直接输出可交付的办公
1611 1
|
5天前
|
SQL 人工智能 前端开发
QoderWake 1.0 正式发布:从桌面里的 Agent,到工作现场的数字员工
QoderWake v1.0正式发布:企业级数字员工团队平台。支持“一句话建岗”,预置10类特训岗位;Waker常驻钉钉/飞书群,@即响应、自动协作、跨任务记忆;具备定时/事件/API多触发方式与统一任务看板;已沉淀27.6万条记忆、12.3万项技能,助力组织实现人机协同增效。
710 1
|
17天前
|
人工智能 自然语言处理 安全
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
本文聚焦阿里云2026年推出的三款自研AI办公产品,清晰拆解千问办公、Qoder Teams、Qoder CN的差异化定位与能力边界:千问办公主打职场全场景提效,支持自然语言指令一键完成PPT生成、数据分析等高频办公任务;Qoder Teams面向程序员团队,深度整合AI代码生成、团队协同与企业知识库能力;Qoder CN则专为金融、政务等强合规场景打造,实现数据不出境与VPC私有化部署。文章同步给出分场景选型指南与最新活动定价,帮助不同类型的企业按需组合产品,实现业务岗、研发岗与强合规场景的AI能力全覆盖。
3868 5
阿里云千问办公、Qoder Teams、Qoder CN区别与选择指南:模型能力、适用场景与最新活动参考
|
8天前
|
人工智能 自然语言处理 安全
阿里云AI数智鉴密:AI 生成内容如何拿到一张"防篡改的身份证"
隐形水印 + C2PA签名:让AI生成内容“持证上岗”。
1142 0
|
9天前
|
网络协议 Linux iOS开发
【2026实测】Wireshark下载+安装+汉化+使用教程(图文版,巨详细)
Wireshark 是一款免费开源的网络协议分析工具,可实时捕获、解析并可视化数据包,助你诊断网络故障、分析通信协议(如HTTP、DNS、TCP等)。支持Windows/macOS/Linux,含中文界面,新手入门便捷。(239字)
|
3天前
|
缓存 测试技术 API
DeepSeek V4.1 Flash 内测接入:改个模型名即可调用(附代码)
DeepSeek V4.1 Flash 内测不用申请,base_url 不变、改个模型名就能调,9/10 到期。本文讲清接入、计费限流与多模态注意点。
685 0
DeepSeek V4.1 Flash 内测接入:改个模型名即可调用(附代码)
|
10天前
|
缓存 数据可视化 开发工具
DeepSeek Harness 怎么更新?dsh 更新完整指南:更新本体(npx、npm、源码)与更新插件两种方式
DeepSeek Harness 的更新分两层:本体更新(npx 自动最新、npm update -g、源码 git pull)与插件更新(插件市场点更新、命令行覆盖安装)。本文按「准备 → 更新本体 → 更新插件 → 更新后检查」四步走,覆盖新手常见疑问。
1250 1
DeepSeek Harness 怎么更新?dsh 更新完整指南:更新本体(npx、npm、源码)与更新插件两种方式