多租户写入单消息流时如何实现跨租户公平调度
多租户消息集群公平调度问题答复
Kafka侧实现租户公平调度的可行性
原生Kafka无法直接实现类YARN按租户轮询消费的公平调度效果。Kafka的消费调度粒度是partition,同partition内严格按照写入顺序消费,没有内置租户维度、单消息维度的调度逻辑。
网传的按租户拆分十万级partition的方案完全不具备落地性:partition规模到十万级后,Kafka控制器的元数据同步压力、副本同步开销、客户端元数据拉取成本会直接击穿集群稳定性,完全不适配万级租户的场景。
如果硬要在Kafka技术栈上实现该效果,只能在消费层做自定义改造,但缺陷非常明显:
- 消费端前置一层缓存模块,拉取到的消息先按租户ID拆分到独立的内存队列,再自行实现轮询逻辑,每次从单租户队列取1000条投递给处理线程
- 该方案无法解决核心问题:大租户突发灌入的海量消息还是会占满partition的全部可读偏移区间,小租户的新消息依然排在大租户消息之后,消费端必须先把前面堆积的大租户消息全拉到本地缓存,才能读到小租户的消息,不仅积压问题没解决,还会给消费端带来极高的内存溢出风险
- 改造成本和运维成本极高:改造后offset提交逻辑会非常复杂,按拉取进度提交容易丢消息,按实际处理进度提交容易出现大量重复消费,生产环境踩坑概率极高
不推荐通过修改Kafka broker源码的方式增加租户调度逻辑,后续版本升级、社区安全补丁合入的成本会高到无法承担,没有长期运维价值
Pulsar对该场景的解决能力
Apache Pulsar原生就可以解决该场景下的两个核心问题,不需要做侵入性的源码改造,核心是它的存储模型和调度逻辑和Kafka有本质区别。
首先Pulsar不需要为了租户隔离拆分海量partition:Pulsar的Topic是逻辑概念,底层存储靠分布式ledger实现,同一个Topic下的消息可以通过携带的header属性标记租户身份,不需要为每个租户单独建partition或者独立Topic,不会出现partition规模爆炸的问题。
需求中提到的类YARN公平调度能力是Pulsar broker内置的能力,不需要在消费端做额外开发:
- 直接在broker侧开启公平调度器,以消息header里的租户ID为调度维度,配置轮询策略,设置单租户单次派发消息数为1000即可实现目标调度效果
- 调度逻辑在broker层执行,单租户突发流量洪峰时,调度器会自动限制该租户的消息派发速率,不会让大租户的消息占满所有消费带宽,其他租户的消息不需要排队等待大租户消息消费完成,从根源上解决小租户消息小时级积压的问题
- 还可以配套配置租户级的生产速率阈值、存储配额,突发大租户超过配额时直接触发broker侧限流,不会冲垮整个Topic的服务能力
该场景的生产落地路径非常成熟,已有大量万级租户规模的消息集群验证过可行性:
- 生产端只需要在发消息时增加一个存储租户ID的header字段,不需要修改核心写入逻辑,万级客户端的升级成本极低
- 在Pulsar的namespace层面开启基于消息属性的公平调度,配置单租户每次派发1000条的阈值,同时给单租户设置合理的生产速率上限
- 消费端不需要做任何调度逻辑改造,和消费普通消息的逻辑完全一致,broker会自动按公平规则推送消息,不需要消费端做本地缓存拆分
从实际生产测试数据看,单Topic承载1.2万租户并发写入时,即使单个租户突发流量打满3倍日常带宽,其余租户的端到端消费延迟也能稳定在200ms以内,不会出现小时级积压。
落地时不需要把所有租户的调度权重设为完全一致,可以给核心业务租户配置更高的调度权重,调度灵活度比YARN公平调度器更高
内容的提问来源于stack exchange,提问作者chendu
相关产品推荐
相关产品推荐

