消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载本文基于 Apache Pulsar 仓库中的 PIP-36Max Message Size提案文档展开介绍 Pulsar 如何把原先硬编码的 5MB 消息上限改造为可通过 broker 配置调整、并通过 Wire Protocol 在连接建立时动态下发给客户端的机制。读完本文你将理解CommandConnected协议字段的演进、maxMessageSize在 Broker/Client/Proxy 三端的完整传递链路以及LengthFieldBasedFrameDecoder动态替换等底层实现细节。1. 背景与动机硬编码消息上限的问题PIP-36 要解决的问题很直接在 Pulsar 早期实现中MaxMessageSize被硬编码在协议层服务端配置无法修改它。用户如果想传输超过默认大小的消息例如较大的二进制对象、批量数据没有任何合法的配置手段。根据 pip-36.md 的 Motivation 部分提案给出两点核心诉求在broker.conf中增加maxMessageSize配置项允许运维人员调整 Broker 可接收的最大消息体由于是 Broker 决定允许接收多大的消息客户端必须知道这个上限才能避免发送失败因此提案在协议中新增max_message_size字段在客户端连接完成时告知其可用的最大消息尺寸。在 Pulsar 中一个消息并不只是 payload它还包括 MessageMetadata、Key 与 Properties、以及 Broker Entry Metadatamagic number、ledger id、entry id、timestamp 等。因此协议层定义的最大消息尺寸需要为这些元数据预留空间这正是后文MESSAGE_SIZE_FRAME_PADDING存在的意义。2. 协议变更CommandConnected 新增 max_message_size 字段PIP-36 对 Wire Protocol 的改动集中在CommandConnected消息上即在 Broker 响应连接CONNECTED时携带最大消息尺寸。提案中的原始定义如下message CommandConnected { required string server_version 1; optional int32 protocol_version 2 [default 0]; // used for telling clients what is the max message size it can use optional int64 max_message_size 3 [default 5 * 1024 * 1024]; }对照当前仓库中实际落地的协议文件 PulsarApi.proto可以看到该字段已存在并沿用了提案中的 tag 编号 3message CommandConnected { required string server_version 1; optional int32 protocol_version 2 [default 0]; optional int32 max_message_size 3; optional FeatureFlags feature_flags 4; }从源码结构看有两个值得注意的落地细节最终实现中max_message_size被定义为int32而非提案中写的int64——由于消息尺寸以字节计int32 已可表示超过 2GB 的尺寸对消息场景足够新增的feature_flagstag 4是后来为 Topic Watchers、Scalable Topics 等能力协商而扩展的说明CommandConnected已成为 Broker 向客户端广播能力与限制的统一通道。协议中的optional语义保证了向后兼容旧版本客户端不会解析该字段旧版本 Broker 不下发该字段时客户端也能按默认值工作。3. 三个核心常量及其语义PIP-36 的 Implement 部分定义了在Commands类中引入的三个与消息尺寸相关的值。当前仓库 Commands.java 中完全一致// default message size for transfer public static final int DEFAULT_MAX_MESSAGE_SIZE 5 * 1024 * 1024; public static final int MESSAGE_SIZE_FRAME_PADDING 10 * 1024; public static final int INVALID_MAX_MESSAGE_SIZE -1;三者语义结合 PIP 原文与源码注释常量值作用DEFAULT_MAX_MESSAGE_SIZE5 MB未显式指定 max message size 时的缺省值同时是 Broker 配置maxMessageSize的默认值MESSAGE_SIZE_FRAME_PADDING10 KB消息元数据MessageMetadata 等预留的空间Netty 帧解码器的最大帧长需要按消息尺寸 padding计算INVALID_MAX_MESSAGE_SIZE-1表示本次 CONNECTED 命令中不携带消息尺寸字段INVALID_MAX_MESSAGE_SIZE的处理逻辑可在 Commands.newConnectedCommand 中直接看到BaseCommand cmd localCmd(Type.CONNECTED); CommandConnected connected cmd.setConnected() .setServerVersion(Pulsar Server PulsarVersion.getVersion()); if (INVALID_MAX_MESSAGE_SIZE ! maxMessageSize) { connected.setMaxMessageSize(maxMessageSize); } ...只有当传入值不是 -1 时才会调用setMaxMessageSize从而让该可选字段保持缺省——这与 PIP 中有时 CONNECTED 消息里不需要携带尺寸可选用 INVALID 值使其不出现在命令中的设计说明完全对应。4. Broker 侧配置项与连接完成时的下发4.1 broker.conf 中的 maxMessageSizePIP-36 提出在broker.conf中添加MaxMessageSize配置。在当前仓库中这一配置项落在 Broker 的服务配置类里ServiceConfiguration.java 定义了private int maxMessageSize Commands.DEFAULT_MAX_MESSAGE_SIZE;默认值直接复用协议常量DEFAULT_MAX_MESSAGE_SIZE5MB保证配置层默认值与协议层默认值天然一致。该类中还有多处引用说明该配置向下传导的下游影响例如 BookKeeper 客户端最大帧长需要按maxMessageSize padding计算、默认 broker maxMessageSize 下 BookKeeper ledger 存储的写入行为等——这解释了为什么 padding 必须参与底层存储链路的容量规划而不仅是传输层的事。4.2 completeConnect把上限写进 CONNECTED 命令PIP 中给出的服务端伪码是// complete the connect and sent newConnected command private void completeConnect(int clientProtoVersion, String clientVersion, long maxMessageSize) { ctx.writeAndFlush(Commands.newConnected(clientProtoVersion, maxMessageSize)); state State.Connected; ... }当前仓库的实际实现位于 ServerCnx.java在通过认证及可选的授权校验后执行// complete the connect and sent newConnected command private void completeConnect(int clientProtoVersion, String clientVersion) { if (service.isAuthenticationEnabled()) { ... maybeScheduleAuthenticationCredentialsRefresh(); } writeAndFlush(Commands.newConnected(clientProtoVersion, maxMessageSize, enableTopicListWatcher, scalableTopicsEnabled, scalableTopicsEnabled service.getPulsar().getConfig().isTransactionCoordinatorScalableTopicsEnabled())); state State.Connected; service.getPulsarStats().recordConnectionCreateSuccess(); ... }可以看到与 PIP 提出的骨架一致maxMessageSize作为completeConnect的调用方持有字段来源于 Broker 配置被写入Commands.newConnected(...)在认证完成后一次性随 CONNECTED 帧下发。从源码结构看这里还顺带下发了 Topic Watchers 与 Scalable Topics 等能力开关——CommandConnected事实上承担了连接期能力握手的角色。5. Client 侧动态替换帧解码器并校验发送尺寸PIP 对客户端行为的要求是客户端不应自行在配置中固定 max message size而应在连接时从 Broker 的 CONNECTED 中获取其支持的上限获取到新尺寸后客户端应替换连接上的LengthFieldBasedFrameDecoder发送消息时与上限比较超过则抛出异常。在仓库中客户端相关实现分布在pulsar-client模块ClientCnx.java 持有与maxMessageSize相关的成员在收到 CONNECTED 后记录 Broker 声明的上限用于后续的发送校验ConnectionHandler.java 与 PulsarChannelInitializer.java 参与连接初始化其中涉及DEFAULT_MAX_MESSAGE_SIZE/MESSAGE_SIZE_FRAME_PADDING用于在拿到 Broker 下发值之前构建初始帧解码器并在 CONNECTED 到达后按新值重建。帧解码器的容量计算逻辑在 FrameDecoderUtil.java 中集中体现maxMessageSize Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4);即 NettyLengthFieldBasedFrameDecoder的最大帧长被设置为maxMessageSize 10KBpadding 正是为消息元数据预留的空间。这就形成了一个闭环Broker 声明的上限 → 客户端解码器容量同步调整 → 双向Broker 发大消息 / 客户端发大消息都在同一容量约束下工作。发送侧的校验行为则对应 PIP 中消息尺寸超过时抛异常的要求客户端在send路径上比较待发送消息的 metadata payload 尺寸与连接级上限超限即抛出MessageTooBigException类错误使失败在本地尽早暴露而不是等 Broker 拒收。6. Proxy 侧的两种连接形态PIP-36 特别讨论了 Proxy 场景其关键在于客户端连接 Proxy 存在两种操作形态。查表lookup阶段客户端只是通过 Proxy 查询目标 Broker 的地址。此时 Proxy 与客户端之间除连接消息外不传输消息数据因此 Proxy 在该连接上无需设置消息尺寸直连转发direct proxy阶段客户端已确定要连接的 BrokerProxy 会为这条连接创建一个 direct proxy。此时 Proxy 与客户端之间会真正转发消息二者需要交换 max message size并在 Proxy-Client 链路上替换LengthFieldBasedFrameDecoder。这一设计的意义在于lookup 连接的帧解码器可以保持轻量默认容量只有在承载真实消息流的直连转发连接上才按协商出的上限调整避免为纯控制面连接浪费解码缓冲。7. 小结从硬编码到可配置的完整链路综合 PIP-36 与当前仓库的实现maxMessageSize的完整链路为配置运维在 broker 侧调整maxMessageSize默认 5MB见 ServiceConfiguration.java协议Broker 在CommandConnected.max_message_sizetag 3见 PulsarApi.proto中携带该值INVALID_MAX_MESSAGE_SIZE时省略该字段BrokerServerCnx.completeConnect在认证完成后通过Commands.newConnected(...)下发见 ServerCnx.javaClient收到 CONNECTED 后更新连接级上限、按maxMessageSize MESSAGE_SIZE_FRAME_PADDING重建LengthFieldBasedFrameDecoder并在发送时本地校验超限即抛异常Proxylookup 连接不设置尺寸直连转发连接则与客户端协商并同步替换解码器。该机制以最小的协议代价一个 optional 字段解决了上限不可配置的问题同时通过optional语义与协议版本协商保持了新旧客户端、新旧 Broker 之间的兼容是理解 Pulsar Wire Protocol 能力协商机制的一个典型样本。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐Apache Pulsar 短 Topic 名Short Topic Names机制详解PIP-11 从设计到实现Apache Pulsar 短 Topic 名Short Topic Names机制详解PIP 11 从设计到实现 导读 Pulsar 从设计之初就是一个消息队列流处理后端微服务消息路由Apache Pulsar 无状态代理Pulsar Proxy架构解析基于 PIP-1 的二进制协议代理设计与实现Apache Pulsar 无状态代理Pulsar Proxy架构解析基于 PIP 1 的二进制协议代理设计与实现 导读 本文以 Apache Pulsa消息队列流处理后端微服务消息路由Apache Pulsar PIP-137 详解基于 Exclusive Producer 的 Pulsar Client 共享状态 API 设计Apache Pulsar PIP 137 详解基于 Exclusive Producer 的 Pulsar Client 共享状态 API 设计 PIP 1消息队列流处理后端微服务消息路由上一篇3分钟终极指南如何让Windows 10/11完美显示iPhone HEIC照片缩略图下一篇HumanEval-Infilling 代码中间填充FIM评测指南lm-evaluation-harness 中的实现与用法创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?