下一篇【第59篇】为什么官方默认只提供gRPC通信——SkyWalking通信设计决策的幽默剖析上一篇【第61篇】数据上报通信扩展——用Kafka传输Trace数据的完整实战一、Agent启动第一步——注册你有没有想过当你启动一个加了SkyWalking Agent的应用时Agent做的第一件事是什么答案出人意料不是上报Trace数据而是注册。Agent需要先告诉OAP嘿我是order-service的实例order-service-pod-001我上线了请记住我。只有注册成功后OAP才知道有这么一个服务实例存在后续的Trace数据才能正确关联。如果这一步失败了后续所有的Trace上报都白搭——OAP不认识这个实例会直接丢弃数据。------------------------------------------------------------------ Agent注册流程标准gRPC版本 | ------------------------------------------------------------------ | | | Agent启动 | | │ | | ↓ | | ┌─────────────┐ | | │ 读取配置文件 │ ← agent.config 中的 agent.service_name | | │ 构建ServiceInstance │ agent.instance_name | | └──────┬───────┘ | | │ | | ↓ | | ┌─────────────────┐ | | │ gRPC Channel建立 │ ← 连接 OAP (默认 127.0.0.1:11800) | | └──────┬──────────┘ | | │ | | ↓ | | ┌─────────────────┐ | | │ ServiceRegister │ ← 调用 RegisterService | | │ (服务注册) │ OAP返回 serviceId | | └──────┬──────────┘ | | │ | | ↓ | | ┌─────────────────┐ | | │ InstanceRegister │ ← 调用 RegisterInstance | | │ (实例注册) │ OAP返回 instanceId | | └──────┬──────────┘ | | │ | | ↓ | | ┌─────────────────┐ | | │ 心跳维持 │ ← 每30秒发送一次心跳 | | │ (keepAlive) │ OAP确认实例存活 | | └─────────────────┘ | | │ | | ↓ | | ┌─────────────────┐ | | │ 正式开始数据上报 │ ← 注册成功后才能做 | | └─────────────────┘ | | | ------------------------------------------------------------------二、gRPC注册实现参考要实现一个HTTP注册扩展我们需要先理解gRPC版本的实现。我们来读一读核心代码2.1 注册服务的protobuf定义// Register.proto (简化版) service Register { // 服务注册 rpc doServiceRegister (Services) returns (ServiceRegisterMapping) {} // 实例注册 rpc doServiceInstanceRegister (ServiceInstances) returns (ServiceInstanceRegisterMapping) {} // 心跳 rpc doHeartbeat (ServiceInstancePingPkg) returns (Commands) {} }2.2 gRPC注册客户端实现简化的关键代码publicclassGRPCChannelManagerimplementsChannelManager,Runnable{privatevolatileGRPCChannelmanagedChannelnull;privatevolatilebooleanreconnecttrue;// 建立gRPC连接Overridepublicvoidrun(){ListStringserversServiceDiscovery.INSTANCE.discover(Config.Collector.BACKEND_SERVICE);for(Stringserver:servers){try{String[]hostPortserver.split(:);managedChannelGRPCChannel.newBuilder(hostPort[0],Integer.parseInt(hostPort[1])).addManagedChannelInterceptor(newAuthenticationInterceptor()).connect();// 连接成功后通知监听器notify(GRPCChannelStatus.CONNECTED);break;}catch(Exceptione){// 重试}}}OverridepublicGRPCChannelgetChannel(){returnmanagedChannel;}}// 服务注册客户端publicclassServiceAndEndpointRegisterClientimplementsGRPCChannelListener{privatevolatileRegisterServiceGrpc.RegisterServiceBlockingStubstub;OverridepublicvoidstatusChanged(GRPCChannelStatusstatus){if(statusGRPCChannelStatus.CONNECTED){// 连接成功后创建stubstubRegisterServiceGrpc.newBlockingStub(GRPCChannelManager.INSTANCE.getChannel());// 开始注册shouldTryRegistertrue;}}// 注册服务privatevoidregisterService(){if(stub!null){ServicesservicesServices.newBuilder().addServices(Service.newBuilder().setServiceName(Config.Agent.SERVICE_NAME)).build();ServiceRegisterMappingmappingstub.doServiceRegister(services);// 保存OAP返回的serviceIdserviceIdmapping.getServices(0).getValue();}}}三、SPI机制的扩展点理解了gRPC版本后问题来了怎么让Agent换用HTTP答案藏在SkyWalking的SPI机制里。SkyWalking使用自定义的SPIServiceLoader模式的变体来管理扩展点。3.1 找到正确的SPI文件------------------------------------------------------------------ | SkyWalking Agent SPI扩展点位置 | ------------------------------------------------------------------ | | | Agent JAR解压后在以下位置可以找到SPI配置文件 | | | | skywalking-agent.jar!/ | | └── META-INF/ | | └── services/ | | ├── org.apache.skywalking.apm.agent.core.remote. | | │ GRPCChannelManager ← 这是我们要替换的扩展点 | | │ | | ├── org.apache.skywalking.apm.agent.core.boot. | | │ BootService ← BootService的SPI声明 | | │ | | └── ... 更多扩展点 ... | | | | 方案A替换GRPCChannelManager整体替换 | | 方案B实现一个新的Reporter只替换数据上报 | | 方案C在Agent外层加代理最不推荐 | | | ------------------------------------------------------------------3.2 扩展点的实现框架// Extension定义SkyWalking自定义的SPI注解Target(ElementType.TYPE)Retention(RetentionPolicy.RUNTIME)publicinterfaceExtension{Stringvalue();// 实现名用于区分不同实现}// 使用示例Extension(grpc)publicclassGRPCChannelManagerimplementsChannelManager,Runnable{// ...}Extension(http)publicclassHTTPChannelManagerimplementsChannelManager,Runnable{// ... 我们的自定义实现}四、实现HTTP注册通信扩展现在开始实现。核心思路将gRPC Channel的逻辑替换为HTTP Client。4.1 HTTPChannelManager实现packageorg.apache.skywalking.apm.agent.core.remote;importio.netty.handler.codec.http.*;importorg.apache.skywalking.apm.agent.core.boot.*;importorg.apache.skywalking.apm.agent.core.conf.*;importorg.apache.skywalking.apm.agent.core.logging.api.*;importjava.io.*;importjava.net.*;importjava.util.*;importjava.util.concurrent.*;/** * HTTP版本的Channel Manager * 替换默认的gRPC Channel */Extension(http)publicclassHTTPChannelManagerimplementsChannelManager,Runnable{privatestaticfinalILogLOGGERLogManager.getLogger(HTTPChannelManager.class);privatevolatilebooleanconnectedfalse;privatevolatileStringbaseUrl;privateListGRPCChannelListenerlistenersnewCopyOnWriteArrayList();privateScheduledExecutorServiceheartbeatExecutor;// HTTP客户端使用Java内置的HttpURLConnectionprivateHttpURLConnectiongetConnection(Stringpath)throwsIOException{URLurlnewURL(baseUrlpath);HttpURLConnectionconn(HttpURLConnection)url.openConnection();conn.setRequestMethod(POST);conn.setRequestProperty(Content-Type,application/json);conn.setRequestProperty(Accept,application/json);conn.setDoOutput(true);conn.setConnectTimeout(5000);conn.setReadTimeout(10000);returnconn;}Overridepublicvoidrun(){// 从配置中获取OAP Server地址StringbackendServiceConfig.Collector.BACKEND_SERVICE;String[]serversbackendService.split(,);for(Stringserver:servers){String[]partsserver.trim().split(:);Stringhostparts[0];Stringportparts.length1?parts[1]:12800;// HTTP端口// 构建基础URLbaseUrlhttp://host:port;// 测试连接if(testConnect()){connectedtrue;// 通知所有监听器连接成功for(GRPCChannelListenerlistener:listeners){listener.statusChanged(GRPCChannelStatus.CONNECTED);}// 启动心跳startHeartbeat();break;}}}privatebooleantestConnect(){try{HttpURLConnectionconngetConnection(/health);intcodeconn.getResponseCode();returncode200;}catch(Exceptione){LOGGER.warn(HTTP connection test failed: {},e.getMessage());returnfalse;}}privatevoidstartHeartbeat(){heartbeatExecutorExecutors.newSingleThreadScheduledExecutor();heartbeatExecutor.scheduleAtFixedRate(this::sendHeartbeat,30,30,TimeUnit.SECONDS);}privatevoidsendHeartbeat(){try{// 构建心跳请求StringheartbeatJsonbuildHeartbeatPayload();HttpURLConnectionconngetConnection(/register/heartbeat);try(OutputStreamosconn.getOutputStream()){os.write(heartbeatJson.getBytes());os.flush();}intcodeconn.getResponseCode();if(code!200){LOGGER.warn(Heartbeat failed with code: {},code);// 解析Commands可能有配置下发}else{// 读取CommandsStringresponsereadResponse(conn);CommandscommandsparseCommands(response);if(commands!null){CommandExecutorService.execute(commands);}}}catch(Exceptione){LOGGER.error(Heartbeat error: {},e.getMessage());}}// 注册服务publicintregisterService(StringserviceName)throwsIOException{StringjsonString.format({\services\:[{\serviceName\:\%s\}]},serviceName);HttpURLConnectionconngetConnection(/register/service);try(OutputStreamosconn.getOutputStream()){os.write(json.getBytes());os.flush();}StringresponsereadResponse(conn);// 解析返回的serviceIdJsonObjectresultparseJson(response);returnresult.getAsJsonArray(services).get(0).getAsJsonObject().get(value).getAsInt();}// 注册实例publicintregisterInstance(intserviceId,StringinstanceName)throwsIOException{StringjsonString.format({\instances\:[{\serviceId\:%d,\instanceName\:\%s\,\time\:%d}]},serviceId,instanceName,System.currentTimeMillis());HttpURLConnectionconngetConnection(/register/instance);try(OutputStreamosconn.getOutputStream()){os.write(json.getBytes());os.flush();}StringresponsereadResponse(conn);JsonObjectresultparseJson(response);returnresult.getAsJsonArray(instances).get(0).getAsJsonObject().get(value).getAsInt();}// 添加监听器publicvoidaddChannelListener(GRPCChannelListenerlistener){listeners.add(listener);}publicbooleanisConnected(){returnconnected;}privateStringreadResponse(HttpURLConnectionconn)throwsIOException{try(BufferedReaderreadernewBufferedReader(newInputStreamReader(conn.getInputStream()))){StringBuildersbnewStringBuilder();Stringline;while((linereader.readLine())!null){sb.append(line);}returnsb.toString();}}privateStringbuildHeartbeatPayload(){returnString.format({\serviceInstanceId\:%d,\time\:%d},getInstanceId(),System.currentTimeMillis());}}4.2 SPI配置文件在扩展Jar包的META-INF/services/目录下创建# org.apache.skywalking.apm.agent.core.remote.ChannelManager httporg.apache.skywalking.apm.agent.core.remote.HTTPChannelManager grpcorg.apache.skywalking.apm.agent.core.remote.GRPCChannelManager4.3 配置切换# agent/config/agent.config# 指定使用HTTP扩展agent.channel_managerhttp agent.collector.backend_serviceoap-server.example.com:12800五、扩展打包与部署5.1 Maven打包配置!-- pom.xml --projectgroupIdcom.example/groupIdartifactIdskywalking-http-channel-extension/artifactIdversion1.0.0/versiondependencies!-- 依赖SkyWalking Agent核心API --dependencygroupIdorg.apache.skywalking/groupIdartifactIdapm-agent-core/artifactIdversion8.16.0/versionscopeprovided/scope/dependency/dependenciesbuildpluginsplugingroupIdorg.apache.maven.plugins/groupIdartifactIdmaven-shade-plugin/artifactIdexecutionsexecutionphasepackage/phasegoalsgoalshade/goal/goals/execution/executions/plugin/plugins/build/project5.2 部署步骤# 1. 编译扩展mvn clean package# 2. 将Jar包放入optional-plugins目录cptarget/skywalking-http-channel-extension-1.0.0.jar\skywalking-agent/optional-plugins/# 3. 配置agent.config# agent.channel_managerhttp# agent.collector.backend_serviceoap-server:12800# 4. 启动应用验证java-javaagent:skywalking-agent.jar-jarmyapp.jar六、测试与验证TestpublicvoidtestHTTPRegistration()throwsException{// 构造配置Config.Collector.BACKEND_SERVICE127.0.0.1:12800;// 创建HTTP Channel ManagerHTTPChannelManagermanagernewHTTPChannelManager();manager.run();// 建立连接assertTrue(Connection should be established,manager.isConnected());// 测试服务注册intserviceIdmanager.registerService(test-service);assertTrue(Service ID should be positive,serviceId0);// 测试实例注册intinstanceIdmanager.registerInstance(serviceId,test-instance);assertTrue(Instance ID should be positive,instanceId0);}七、性能与兼容性注意事项------------------------------------------------------------------ | HTTP vs gRPC注册的性能特征对比 | ------------------------------------------------------------------ | | | 指标 gRPC HTTP | | ─────────────────────────────────────────────────────────────── │ | 首次注册耗时 快速长连接 中等需建立TCP连接 | | 心跳开销 极低连接复用 中等每次需HTTP请求/响应 | | 网络开销 小protobuf 大JSON文本 | | 防火墙友好度 差需独立端口 好可用80/443端口 | | 代理支持 差 好 | | 双向通信 原生支持 需轮询实现 | | | ------------------------------------------------------------------八、总结本文展示了SkyWalking通信扩展的完整实现思路。核心要点注册是Agent启动的第一道门槛注册失败意味着后续所有监控数据都无法上报SPI机制是扩展的入口通过注解Extension可以注册自定义实现HTTP作为替代方案适合网络受限环境但性能不如gRPC扩展Jar放optional-plugins目录即可生效下一篇我们深入更具实用价值的扩展——用Kafka传输Trace数据。下一篇【第59篇】为什么官方默认只提供gRPC通信——SkyWalking通信设计决策的幽默剖析上一篇【第61篇】数据上报通信扩展——用Kafka传输Trace数据的完整实战