【SkyWalking从入门到精通】第60篇:探针注册通信扩展——基于HTTP的SPI注册实现完全指南
下一篇【第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{privatevolatileGRPCChannelmanagedChannel=null;privatevolatilebooleanreconnect=true;// 建立gRPC连接@Overridepublicvoidrun(){List<String>servers=ServiceDiscovery.INSTANCE.discover(Config.Collector.BACKEND_SERVICE);for(Stringserver:servers){try{String[]hostPort=server.split(":");managedChannel=GRPCChannel.newBuilder(hostPort[0],Integer.parseInt(hostPort[1])).addManagedChannelInterceptor(newAuthenticationInterceptor()).connect();// 连接成功后通知监听器notify(GRPCChannelStatus.CONNECTED);break;}catch(Exceptione){// 重试}}}@OverridepublicGRPCChannelgetChannel(){returnmanagedChannel;}}// 服务注册客户端publicclassServiceAndEndpointRegisterClientimplementsGRPCChannelListener{privatevolatileRegisterServiceGrpc.RegisterServiceBlockingStubstub;@OverridepublicvoidstatusChanged(GRPCChannelStatusstatus){if(status==GRPCChannelStatus.CONNECTED){// 连接成功后创建stubstub=RegisterServiceGrpc.newBlockingStub(GRPCChannelManager.INSTANCE.getChannel());// 开始注册shouldTryRegister=true;}}// 注册服务privatevoidregisterService(){if(stub!=null){Servicesservices=Services.newBuilder().addServices(Service.newBuilder().setServiceName(Config.Agent.SERVICE_NAME)).build();ServiceRegisterMappingmapping=stub.doServiceRegister(services);// 保存OAP返回的serviceIdserviceId=mapping.getServices(0).getValue();}}}三、SPI机制的扩展点
理解了gRPC版本后,问题来了:怎么让Agent换用HTTP?
答案藏在SkyWalking的SPI机制里。SkyWalking使用自定义的SPI(ServiceLoader模式的变体)来管理扩展点。
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)public@interfaceExtension{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{privatestaticfinalILogLOGGER=LogManager.getLogger(HTTPChannelManager.class);privatevolatilebooleanconnected=false;privatevolatileStringbaseUrl;privateList<GRPCChannelListener>listeners=newCopyOnWriteArrayList<>();privateScheduledExecutorServiceheartbeatExecutor;// HTTP客户端(使用Java内置的HttpURLConnection)privateHttpURLConnectiongetConnection(Stringpath)throwsIOException{URLurl=newURL(baseUrl+path);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地址StringbackendService=Config.Collector.BACKEND_SERVICE;String[]servers=backendService.split(",");for(Stringserver:servers){String[]parts=server.trim().split(":");Stringhost=parts[0];Stringport=parts.length>1?parts[1]:"12800";// HTTP端口// 构建基础URLbaseUrl="http://"+host+":"+port;// 测试连接if(testConnect()){connected=true;// 通知所有监听器连接成功for(GRPCChannelListenerlistener:listeners){listener.statusChanged(GRPCChannelStatus.CONNECTED);}// 启动心跳startHeartbeat();break;}}}privatebooleantestConnect(){try{HttpURLConnectionconn=getConnection("/health");intcode=conn.getResponseCode();returncode==200;}catch(Exceptione){LOGGER.warn("HTTP connection test failed: {}",e.getMessage());returnfalse;}}privatevoidstartHeartbeat(){heartbeatExecutor=Executors.newSingleThreadScheduledExecutor();heartbeatExecutor.scheduleAtFixedRate(this::sendHeartbeat,30,30,TimeUnit.SECONDS);}privatevoidsendHeartbeat(){try{// 构建心跳请求StringheartbeatJson=buildHeartbeatPayload();HttpURLConnectionconn=getConnection("/register/heartbeat");try(OutputStreamos=conn.getOutputStream()){os.write(heartbeatJson.getBytes());os.flush();}intcode=conn.getResponseCode();if(code!=200){LOGGER.warn("Heartbeat failed with code: {}",code);// 解析Commands(可能有配置下发)}else{// 读取CommandsStringresponse=readResponse(conn);Commandscommands=parseCommands(response);if(commands!=null){CommandExecutorService.execute(commands);}}}catch(Exceptione){LOGGER.error("Heartbeat error: {}",e.getMessage());}}// 注册服务publicintregisterService(StringserviceName)throwsIOException{Stringjson=String.format("{\"services\":[{\"serviceName\":\"%s\"}]}",serviceName);HttpURLConnectionconn=getConnection("/register/service");try(OutputStreamos=conn.getOutputStream()){os.write(json.getBytes());os.flush();}Stringresponse=readResponse(conn);// 解析返回的serviceIdJsonObjectresult=parseJson(response);returnresult.getAsJsonArray("services").get(0).getAsJsonObject().get("value").getAsInt();}// 注册实例publicintregisterInstance(intserviceId,StringinstanceName)throwsIOException{Stringjson=String.format("{\"instances\":[{\"serviceId\":%d,"+"\"instanceName\":\"%s\",\"time\":%d}]}",serviceId,instanceName,System.currentTimeMillis());HttpURLConnectionconn=getConnection("/register/instance");try(OutputStreamos=conn.getOutputStream()){os.write(json.getBytes());os.flush();}Stringresponse=readResponse(conn);JsonObjectresult=parseJson(response);returnresult.getAsJsonArray("instances").get(0).getAsJsonObject().get("value").getAsInt();}// 添加监听器publicvoidaddChannelListener(GRPCChannelListenerlistener){listeners.add(listener);}publicbooleanisConnected(){returnconnected;}privateStringreadResponse(HttpURLConnectionconn)throwsIOException{try(BufferedReaderreader=newBufferedReader(newInputStreamReader(conn.getInputStream()))){StringBuildersb=newStringBuilder();Stringline;while((line=reader.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 http=org.apache.skywalking.apm.agent.core.remote.HTTPChannelManager grpc=org.apache.skywalking.apm.agent.core.remote.GRPCChannelManager4.3 配置切换
# agent/config/agent.config# 指定使用HTTP扩展agent.channel_manager=http agent.collector.backend_service=oap-server.example.com:12800五、扩展打包与部署
5.1 Maven打包配置
<!-- pom.xml --><project><groupId>com.example</groupId><artifactId>skywalking-http-channel-extension</artifactId><version>1.0.0</version><dependencies><!-- 依赖SkyWalking Agent核心API --><dependency><groupId>org.apache.skywalking</groupId><artifactId>apm-agent-core</artifactId><version>8.16.0</version><scope>provided</scope></dependency></dependencies><build><plugins><plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-shade-plugin</artifactId><executions><execution><phase>package</phase><goals><goal>shade</goal></goals></execution></executions></plugin></plugins></build></project>5.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_manager=http# agent.collector.backend_service=oap-server:12800# 4. 启动应用验证java-javaagent:skywalking-agent.jar-jarmyapp.jar六、测试与验证
@TestpublicvoidtestHTTPRegistration()throwsException{// 构造配置Config.Collector.BACKEND_SERVICE="127.0.0.1:12800";// 创建HTTP Channel ManagerHTTPChannelManagermanager=newHTTPChannelManager();manager.run();// 建立连接assertTrue("Connection should be established",manager.isConnected());// 测试服务注册intserviceId=manager.registerService("test-service");assertTrue("Service ID should be positive",serviceId>0);// 测试实例注册intinstanceId=manager.registerInstance(serviceId,"test-instance");assertTrue("Instance ID should be positive",instanceId>0);}七、性能与兼容性注意事项
+------------------------------------------------------------------+ | 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数据的完整实战