ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Pinpoint Channel 模块解析:基于 Redis/Kafka 的统一消息通道抽象与 ChannelService 需求-供给通信框架

2026/9/23 1:11:42 拓冰建站 浏览量
Pinpoint Channel 模块解析:基于 Redis/Kafka 的统一消息通道抽象与 ChannelService 需求-供给通信框架 后端可观测性APM链路追踪微服务【免费下载链接】pinpointAPM, (Application Performance Management) tool for large-scale distributed systems.项目地址https://gitcode.com/gh_mirrors/pi/pinpoint点击查看免费下载导读本文围绕 Pinpoint 仓库中 channel/README.md 讲解的channel 模块展开该模块将“从一个系统部件向另一个系统部件发送消息、并向未知数量的订阅者广播事件”的能力抽象为统一的Channel通道概念底层可以由 Redis、Kafka 等分布式中间件实现。读者读完本文后将掌握如何通过ChannelProviderRepository URI 获取发布/订阅通道、如何完成“发布者-订阅者”Hello World 通信以及如何使用ChannelServiceProtocol搭建类似 RPC 的“需求-供给demand-supply”双向服务含 Mono 单响应与 Flux 多响应两种模式并理解其在 Pinpoint 监控系统中的实际落地形态。一、模块定位把“通道”从中间件细节中抽象出来Pinpoint 是一个面向大规模分布式系统的 APM应用性能管理工具。在它的 Collector、Web、Agent 等多个进程之间存在大量实时消息传递与事件广播需求——例如实时流数据、服务间即时交互。channel模块的核心目标就是抽象出 channel 概念屏蔽底层的 Redis、Kafka 等具体实现。从源码看channel模块位于 channel/src/main/java/com/navercorp/pinpoint/channel核心接口非常精炼Channel.java同时具备发布与订阅能力的通道是PubChannel与SubChannel的组合。其 Javadoc 明确指出在大多数场景下配对的PubChannel与SubChannel位于网络的不同侧由 Redis、Kafka 等分布式系统实现。PubChannel.java仅有一个方法void publish(byte[] content)负责把字节数组发布到与之相连的 SubChannel。SubChannel.java负责注册/注销针对传入字节数组的处理器提供subscribe(SubConsumer)与unsubscribe(Subscription)。SubConsumer.java订阅处理器boolean consume(byte[] content)的返回值用于告知实现方该消息是否被成功消费。Subscription.java订阅句柄调用unsubscribe()可停止监听并回收资源。通道传输的载荷统一为byte[]序列化与反序列化Serde被独立抽象在 serde 包下SerdeT同时继承 Spring 的SerializerT与DeserializerT默认实现为基于 Jackson 的 JacksonSerde.java。二、ChannelRepository按 URI 路由到具体 ChannelProvider2.1 单仓库 多 Provider 的注册模型README 强调了一个关键约定JVM 中只应存在唯一的 channel repository它由若干ChannelProvider与名字scheme的配对创建。示例代码如下ChannelRepository repository new ChannelRepository(List.of( ChannelProviderRegistry.of(redis, new RedisChannelProvider(redis)), ChannelProviderRegistry.of(kafka, new KafkaChannelProvider(kafka)) ))需要说明的是README 中的ChannelRepository在仓库源码中的实际对应类名为 ChannelProviderRepository.java。其构造器接收IterableChannelProviderRegistry在内部以LinkedHashMap按 scheme 索引各 Providerpublic ChannelProviderRepository(IterableChannelProviderRegistry registries) { MapString, ChannelProvider providerMap new LinkedHashMap(); for (ChannelProviderRegistry registry: registries) { providerMap.put(registry.getScheme(), registry.getProvider()); } this.providerMap providerMap; }ChannelProviderRegistry.java一个“scheme → ChannelProvider”的配对静态工厂of(scheme, provider)创建。ChannelProvider.java同时继承PubChannelProvider与SubChannelProvider按 key 提供发布通道和订阅通道并提供静态方法pair(pub, sub)将两个独立 Provider 组合为 PairedChannelProvider.java。2.2 通过 URI 定位通道ChannelProviderRepository的 Javadoc 明确了 URI 格式约定scheme://key e.g. redis://hello-world-topic?param1value1param2value2 e.g. kafka://hello-world-topic?param1value1param2value2即scheme 决定选用哪个 ProviderURI.getSchemeSpecificPart()//之后的部分作为传给 Provider 的 key。其路由逻辑如下public PubChannel getPubChannel(URI uri) { return getChannelProvider(uri).getPubChannel(uri.getSchemeSpecificPart()); } public SubChannel getSubChannel(URI uri) { return getChannelProvider(uri).getSubChannel(uri.getSchemeSpecificPart()); } private ChannelProvider getChannelProvider(URI uri) { ChannelProvider provider this.providerMap.get(uri.getScheme()); if (provider null) { throw new IllegalArgumentException(Scheme not found: uri.getScheme()); } return provider; }核心语义只要通过相同 key 取得的 PubChannel 与 SubChannel 成对出现即便位于不同进程、不同网络节点也能相互通信。若 URI 的 scheme 未注册会抛出IllegalArgumentException(Scheme not found: ...)。2.3 内存实现最直观的通道示例测试目录下的 MemoryChannelProvider.java 是一个纯内存实现完整演示了 Channel 三要素发布、订阅、退订的最小实现内部用MapString, Channel按 key 缓存通道MemChannel用SetSubConsumer维护消费者集合publish时遍历调用每个consumer.consume(content)subscribe返回持有消费者的订阅句柄unsubscribe则从集合移除。它是理解真实 Provider如 Redis Pub/Sub、Redis Stream、Kafka的最佳入门参照。2.4 Spring 装配模块还提供 ChannelSpringConfig.java将ChannelProviderRepository声明为 Spring Bean自动收集所有ChannelProviderRegistryBean并提供基于 JacksonObjectMapper的JsonSerdeFactoryConfiguration(proxyBeanMethods false) public class ChannelSpringConfig { Bean public ChannelProviderRepository channelProviderRepository(ListChannelProviderRegistry registries) { return new ChannelProviderRepository(registries); } Bean public JsonSerdeFactory jsonSerdeFactory(ObjectMapper objectMapper) { return new JacksonSerdeFactory(objectMapper); } }三、Hello World发布者-订阅者广播通信README 给出完整的双实例示例。假设实例 A 是订阅方subscriber实例 B 是发布方publisher双方使用完全相同的 URIredis://system-out?paramfoo订阅方即可收到发布方推送的字节内容。订阅方Instance-subscriberURI uri URI.create(redis://system-out?paramfoo); SubChannel subChannel repository.getSubChannel(uri); subChannel.subscribe(message - { System.out.println(new String(message)); });发布方Instance-publisherURI uri URI.create(redis://system-out?paramfoo); PubChannel pubChannel repository.getPubChannel(uri); pubChannel.publish(Hello, world!.getBytes());运行后订阅方进程应打印Hello, world!。这里paramfoo这样的 query 参数随schemeSpecificPart一并作为 key 传给 Provider可作为区分同一 topic 不同子通道的手段具体是否解析取决于 Provider 实现。这个示例同时体现了 channel 的广播能力channel 可以向未知数量的多个订阅者广播同一事件——每个订阅者通过subscribe注册自己的SubConsumerProvider 将消息投递给所有已注册的消费者参见MemoryChannelProvider.MemChannel.publish的遍历逻辑。四、ChannelService类似 RPC 的需求-供给服务4.1 设计动机与模型除了纯广播通道模块还提供ChannelService实现用于管理系统各部分之间即时的需求-供给demand-supply交互。它与传统 RPC 调用非常相似但区别在于需求demand会被发送给所有正在监听该服务的服务器。工作模型如下ChannelServiceServer必须先于客户端启动并“在网络上”提供服务每台服务器监听为服务预留的需求通道demand channel捕获所有需求收到需求后恰好 0 或 1 台服务器应向供给通道supply channel供给数据0 表示无服务器可响应1 表示由某台服务器响应ChannelServiceClient发出需求并在供给通道上等待结果。因此所有 Server 与 Client 都必须共享一份ChannelServiceProtocol其中包含服务本身、需求通道与供给通道的完整信息。从源码看协议体系位于 service 包ChannelServiceProtocol.java同时继承ChannelServiceServerProtocol与ChannelServiceClientProtocol其 Javadoc 指出按供给数量可将服务分为 Mono 与 Flux 两类——Mono 只发送一个供给Flux 发送多个供给。ChannelServiceProtocolBuilder.java提供buildMono()与buildFlux()两个构建入口。4.2 协议定义完整配置项说明README 给出了ChannelServiceProtocol的完整构建示例Mono 模式结合 ChannelServiceProtocolBuilder.java 源码各配置项含义如下ChannelServiceProtocol protocol ChannelServiceProtocol.String, Longbuilder() .setDemandSerde(JacksonSerde.byClass(objectMapper, String.class)) .setDemandPubChannelURIProvider(demand - URI.create(redis:char-count:demand)) .setDemandSubChannelURI(URI.create(redis:char-count:demand)) .setSupplySerde(JacksonSerde.byClass(objectMapper, Long.class)) .setSupplyChannelURIProvider(demand - URI.create(redis:char-count:supply: demand.hashCode())) .setRequestTimeout(Duration.ofSeconds(3)) .buildMono();配置项类型作用说明setDemandSerdeSerdeD需求对象泛型 D的序列化/反序列化器示例用 Jackson 将String编解码setDemandPubChannelURIProviderFunctionD, URI根据具体需求内容动态生成“发布需求”所用通道的 URIsetDemandSubChannelURIURI所有服务器监听需求所用通道的固定 URI与上面的 Pub URI 配对setSupplySerdeSerdeS供给对象泛型 S的序列化/反序列化器示例中供给是LongsetSupplyChannelURIProviderFunctionD, URI根据需求内容动态计算供给通道 URI用于把响应路由回对应客户端示例按demand.hashCode()区分setRequestTimeoutDuration请求超时时间仅 Mono 模式生效示例为 3 秒setDemandIntervalDuration需求重发/轮询间隔默认Duration.ZEROFlux 模式相关setBufferSizeint缓冲大小默认4Flux 模式相关setChannelStateFnFunctionS, ChannelState从供给结果提取通道状态Flux 模式相关URI scheme 简写说明redis:char-count:demand这种不带//的写法其getScheme()同样是redisgetSchemeSpecificPart()为char-count:demand因此同样能被ChannelProviderRepository正确路由——URI 的格式是scheme://key或scheme:key均可。4.3 服务端Server服务器端通过ChannelServiceServer.buildMono创建并监听ChannelServiceServer.buildMono(repository, protocol, demand - demand.length()).listen();第三个参数是后端处理函数ChannelServiceMonoBackendD, S收到需求后执行计算并返回供给。这里demand - demand.length()表示“统计需求字符串长度”。源码中 ChannelServiceServer.java 同时提供buildMono与buildFlux两个工厂方法二者共用ChannelServiceServerImpl由于该接口继承InitializingBean作为 Spring Bean 注册时其afterPropertiesSet()会自动调用listen()无需手动启动。4.4 客户端Client客户端通过ChannelServiceClient.buildMono创建MonoChannelServiceClient client ChannelServiceClient.buildMono(repository, protocol); Long result client.demand(Hello, World!).block(); // 13demand(Hello, World!)返回MonoLongblock()等待结果为13即Hello, World!的字符数。从 MonoChannelServiceClientImpl.java 可以看到其底层调用链先订阅供给通道subscribe再通过需求发布通道发布序列化后的需求getDemandPubChannel(demand).publish(...)对整个Mono应用timeout(protocol.getRequestTimeout())超时控制当Mono被 dispose例如block()完成或调用方取消时自动退订sink.onDispose(subscription::unsubscribe)。这套“先订阅、后发布、超时保护、自动退订”的时序是理解 ChannelService 客户端实现的关键。对应的服务端与客户端协议接口分别为 ChannelServiceServerProtocol.java 与 ChannelServiceClientProtocol.java。五、Mono 与 Flux 两种服务形态按供给数量ChannelService 分为两类见ChannelServiceProtocol的 Javadoc 与ChannelServiceProtocolBuilder的buildMono()/buildFlux()Mono 服务一次需求对应至多一个供给适合“请求-响应”型交互如查询统计结果。客户端接口为 MonoChannelServiceClient.javaMonoS request(D demand)服务端后端为ChannelServiceMonoBackend需要配置requestTimeout。Flux 服务一次需求对应多个供给数据流适合订阅型交互如持续的实时数据推送。客户端接口为 FluxChannelServiceClient.java服务端后端为ChannelServiceFluxBackend可额外配置demandInterval需求重发间隔、bufferSize默认 4与channelStateFn供给状态判定。测试目录中 MonoChannelServiceTest.java 与 FluxChannelServiceTest.java 分别覆盖两种形态的完整闭环ChannelServiceProtocolTest.java 则验证协议构建与 URI 推导逻辑可作为理解该框架行为的第一手资料。六、真实 Provider 与模块依赖关系6.1 Redis 的三种通道实现channel 模块只负责抽象真实实现由redis模块提供路径 redis/src/main/java/com/navercorp/pinpoint/channel/redis按 Redis 数据模型可分为三类Pub/SubRedisPubChannelProvider.java借助 Redis Pub/Sub 实现即时广播对应配置 RedisPubSubConfig.javaStreamRedisStreamPubChannelProvider.java基于 Redis Stream 实现可提供一定程度的持久化与消费组语义对应配置 RedisStreamConfig.javaKVRedisKVPubChannelProvider.java基于普通键值读写实现对应配置 RedisKVChannelConfig.java。这些 Provider 通过ChannelProviderRegistry.of(redis, provider)注册进唯一的ChannelProviderRepository从而支持redis://...形式的 URI。6.2 在 Pinpoint 中的实际应用场景从仓库结构与依赖关系看channel 抽象被用于 Pinpoint 多进程间的实时交互。例如realtime模块realtime与otlplog/otlpmetric等模块的数据通道以及web与collector之间的即时服务交互均以 channel 为传输底座相关依赖可在各模块 pom.xml 中确认。这正是本模块存在的意义上层业务只面向ChannelProviderRepository URI 编程底层从 Redis 切换到 Kafka 或其它实现时上层代码无需改动只需更换注册的 Provider 与 URI scheme。七、小结与最佳实践要点围绕channel/README.md的核心内容结合仓库源码可以沉淀出以下实践要点JVM 内保持单一ChannelProviderRepository所有通道统一从该仓库按 URI 获取避免多实例造成通道路由混乱URI 即寻址协议scheme决定底层实现redis、kafka等schemeSpecificPart决定通道 key相同 key 的 Pub/Sub 通道在分布式环境下天然互通载荷统一为byte[]业务对象编解码交给Serde默认 Jacksonchannel 层与具体序列化格式解耦广播用 Pub/Sub 通道一个发布者可面向多个未知订阅者广播事件即时供需交互用 ChannelService共享ChannelServiceProtocol后Server 监听需求通道、Client 发布需求并等待供给需注意“Server 必须先于 Client 就绪”且一次需求至多一台服务器响应供给按响应数量选型单响应用buildMono()配requestTimeout多响应流用buildFlux()配demandInterval、bufferSize、channelStateFn直接复用测试用例channel/src/test 下的MemoryChannelProvider与 Mono/Flux 服务测试是理解整个抽象的最短路径也可作为自研 Provider 的参考实现模板。如需继续深入可依次阅读 Channel.java、ChannelProviderRepository.java、ChannelServiceProtocolBuilder.java 与 MonoChannelServiceClientImpl.java即可完整掌握从通道抽象到服务框架的整条链路。赞分享后端可观测性APM链路追踪微服务【免费下载链接】pinpointAPM, (Application Performance Management) tool for large-scale distributed systems.项目地址https://gitcode.com/gh_mirrors/pi/pinpoint点击查看免费下载相关推荐KubeEdge Viaduct 云边通信框架深度解析基于 Protobuf 的 WebSocket / QUIC 消息通道KubeEdge Viaduct 云边通信框架深度解析基于 Protobuf 的 WebSocket / QUIC 消息通道 Viaduct 是 KubeEd云原生边缘计算物联网容器编排边缘网关ABP框架教程模块化CRM系统开发(7) - 基于消息的事件通信ABP框架教程模块化CRM系统开发 7 基于消息的事件通信 前言 在模块化系统设计中模块间的通信是一个关键问题。ABP框架提供了强大的事件总线系统支持模块后端Web框架微服务前端Pintree主题定制教程打造独一无二的个性化导航界面Pintree主题定制教程打造独一无二的个性化导航界面 Pintree是一款能将浏览器书签快速转化为目录网站的开源工具让用户在几分钟内即可构建属于自己的导航人工智能大模型RAG创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考