尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Zipkin ScribeCollector 详解:基于 Scribe 协议接收 Thrift+Base64 编码 Span 的采集器

Zipkin ScribeCollector 详解:基于 Scribe 协议接收 Thrift+Base64 编码 Span 的采集器 可观测性后端微服务【免费下载链接】zipkinZipkin is a distributed tracing system项目地址https://gitcode.com/gh_mirrors/zip/zipkin点击查看免费下载导读本文围绕 zipkin-collector/scribe/README.md 展开系统讲解 Zipkin 中ScribeCollector的功能定位、数据编码格式TBinaryProtocol 大端序 Base64、Builder配置项默认监听 9410 端口、默认消费zipkin分类、Spring Boot 集成参数以及从 Netty 帧解码到 Armeria Thrift 服务、再到异步写入存储的完整调用链。读完本文你将能够在生产或自研链路中接入遗留 Scribe 客户端如 Finagle 生态并正确构造、发送、排障 Scribe 格式的 Zipkin Span。ScribeCollector 是什么ScribeCollector是 Zipkin 众多采集器之一位于 zipkin-collector/scribe 模块其 Maven artifact 为zipkin-collector-scribepom.xml 中模块名为 Collector: Scribe (Legacy)即遗留兼容采集器。它负责在指定的 Scribe 分类category下接收 Scribe 日志并把日志条目中包含的单个 Span 解码后推送到存储后端。它的定位来自 Zipkin 的历史生态Scribe 是 Facebook 开源的日志聚合协议早期 Twitter 的 Finagle RPC 框架通过 Scribe 协议上报 Zipkin 追踪数据。因此ScribeCollector主要服务于这类历史客户端使其采集到的 Span 仍能进入 Zipkin 的现代zipkin2数据管道。其核心类位于 ScribeCollector.javapublic final class ScribeCollector extends CollectorComponent { public static Builder newBuilder() { return new Builder(); } ... }数据编码格式TBinaryProtocol 大端序 Base64README 的Encoding一节给出了整个协议的核心约定这也是接入 ScribeCollector 的客户端必须遵循的格式Scribe message 是一个 TBinaryProtocol 大端序big-endian随后经过 Base64 编码的 Span。Base64 Basic 与 MIME 两种方案均被支持。用伪代码描述为serialized writeTBinaryProtocol(span) encoded base64(serialized) scribe.log(category zipkin, message encoded)即每个LogEntry.message里存放的是Thrift 二进制 Base64 文本分类字段category用于标识这是 Zipkin 的 Span 数据。LogEntry的 Thrift 结构定义可以在生成代码 LogEntry.java 中看到只有两个字段categoryfield 1STRING日志分类默认取值为zipkinmessagefield 2STRINGBase64 编码后的 Span 二进制。解码侧的真实实现在 ScribeSpanConsumer.java 的Log方法中解码流程与 README 的伪代码一一对应byte[] bytes logEntry.message.getBytes(StandardCharsets.ISO_8859_1); bytes Base64.getMimeDecoder().decode(bytes); // finagle-zipkin uses mime encoding byteCount bytes.length; spans.add(SpanBytesDecoder.THRIFT.decodeOne(bytes));可以提炼出三个关键实现事实Base64 解码使用Base64.getMimeDecoder()。MIME 解码器兼容 Basic 与 MIME 字母表这正是 README 所说Base64 Basic 和 MIME 方案都被支持的源码依据注释还特别注明finagle-zipkin 客户端使用 MIME 编码即可能携带换行符等 MIME 特性。Thrift 解码SpanBytesDecoder.THRIFT.decodeOne(bytes)把二进制还原为 zipkin2 的Span对象。这里的THRIFT解码器与zipkin模块中的 SpanBytesDecoder.java 相对应。字节计数解码后的二进制长度被累加到byteCount供采集器指标统计见下文指标。Base64 编码端的示例对应地客户端发送前的编码可用 Java 标准库完成import java.util.Base64; import zipkin2.Span; import zipkin2.codec.SpanBytesEncoder; Span span /* 构造或从 V1Span 转换得到的 zipkin2.Span */; byte[] thriftBytes SpanBytesEncoder.THRIFT.encode(span); String encoded new String(Base64.getEncoder().encode(thriftBytes), StandardCharsets.UTF_8);上述写法与测试 ScribeSpanConsumerTest.java 中构造encodedSpan的方式一致byte[] bytes SpanBytesEncoder.THRIFT.encode(v2); String encodedSpan new String(Base64.getEncoder().encode(bytes), UTF_8);Builder 配置与默认值README 明确指出zipkin2.collector.scribe.ScribeCollector.Builder提供了默认配置监听 9410 端口接收分类为 zipkin 的日志条目。这两个默认值可以直接在 ScribeCollector.java 的Builder字段初始化中确认Collector.Builder delegate Collector.newBuilder(ScribeCollector.class); CollectorMetrics metrics CollectorMetrics.NOOP_METRICS; String category zipkin; int port 9410;Builder提供的方法及说明如下方法说明默认值storage(StorageComponent)设置存储组件解码后的 Span 将异步写入该存储无必须设置metrics(CollectorMetrics)设置指标收集器会调用metrics.forTransport(scribe)为 scribe 传输建立独立指标命名空间CollectorMetrics.NOOP_METRICSsampler(CollectorSampler)设置采样器在写入存储前决定 Span 是否被丢弃由Collector.newBuilder提供默认采样category(String)指定从哪个 Scribe 分类消费 Span传入null会抛出NullPointerExceptionzipkinport(int)指定监听端口9410build()构造ScribeCollector实例—其中port(int)传入0时表示使用系统分配的匿名端口可用于测试场景见下文的测试证据。生命周期与健康检查ScribeCollector还实现了CollectorComponent的生命周期接口start()启动内部的NettyScribeServer如果端口已被占用会抛出RuntimeException(Could not start scribe server.)见 NettyScribeServer.java。check()检查服务是否在运行未启动或通道失活时返回CheckResult.failed(...)。port()在start()之前返回 0启动后返回实际监听端口。close()关闭监听通道并优雅关闭事件循环组。toString()输出形如ScribeCollector{port端口, categoryzipkin}的摘要信息用于健康检查端点与日志展示。这些行为均有单元测试佐证见 ScribeCollectorTest.javacheck_failsWhenNotStarted未启动时check()失败启动后check()通过anonymousPortport(0)时启动前port()为 0启动后为非 0start_failsWhenCantBindPort重复占用同一端口时启动抛出Could not start scribe server.toStringContainsOnlySummaryInformationtoString()只包含端口与分类摘要信息不泄露敏感数据。在 zipkin-server 中启用 Scribe 采集器zipkin-server通过 Spring Boot 自动装配集成 ScribeCollector相关配置类为 ZipkinScribeCollectorConfiguration.java。三个关键配置项配置项含义默认值zipkin.collector.scribe.enabled是否启用 Scribe 采集器只有显式设置为true才加载未启用zipkin.collector.scribe.category消费的 Scribe 分类zipkinzipkin.collector.scribe.port监听端口9410启用方式例如在启动 zipkin-server 时传入 JVM 参数-Dzipkin.collector.scribe.enabledtrue配置类的装配条件为ConditionalOnClass(ScribeCollector.class)ConditionalOnProperty(value zipkin.collector.scribe.enabled, havingValue true)即只有依赖中存在 scribe 采集器 jar 且显式开启时才生效。Bean 使用initMethod start启动时会阻塞直到 scribe 端口开始监听若端口冲突则直接启动失败属于快速失败设计。内部工作原理从 Netty 帧到 Armeria Thrift 服务ScribeCollector的底层接收端由两个类协作完成理解它们有助于排查协议与性能问题。1. NettyScribeServer帧解码与 TCP 服务NettyScribeServer.java 使用 NettyServerBootstrap建立 TCP 服务管道中注册了两个 handlerch.pipeline() .addLast(new LengthFieldBasedFrameDecoder(Integer.MAX_VALUE, 0, 4, 0, 4)) .addLast(new ScribeInboundHandler(scribe));LengthFieldBasedFrameDecoder(Integer.MAX_VALUE, 0, 4, 0, 4)解析的是 Thrift framed transport 的4 字节大端长度前缀length field offset 0、length 4、lengthAdjustment 0、stripBytes 4这与 Scribe/Thrift 的TFramedTransport客户端行为对应——在集成测试 ITScribeCollector.java 中客户端正是使用new TFramedTransport(new TSocket(localhost, server.port()))建立连接后调用Scribe.Client.Log(entries)。事件循环方面bossGroup由 Armeria 的EventLoopGroups.newEventLoopGroup(1)创建单线程 acceptorworkerGroup复用 Armeria 的CommonPools.workerGroup()共享线程池。2. ScribeInboundHandler把 Netty ByteBuf 转为 Armeria Thrift 请求ScribeInboundHandler.java 是连接 Netty 与 Armeria 的桥梁它把收到的原始ByteBuf包装成 Armeria 的HttpRequest请求头固定为POST /internal/zipkin-thriftrpc、Content-Type: application/x-thrift、User-Agent: Zipkin/ScribeInboundHandler通过THttpService.of(scribe)scribe即ScribeSpanConsumer把 Armeria 的 Thrift 服务暴露到这条内部 HTTP 通道上响应被聚合后按照请求顺序维护pendingResponses保证乱序完成的响应按请求序号顺序写回从而与 Scribe 客户端的同步调用语义保持一致异常或通道关闭时释放所有挂起的ByteBuf并关闭连接。也就是说对外呈现的是传统的 Scribe framed TCP 协议对内则复用了 Armeria 的 Thrift 服务栈依赖armeria-thrift0.18见 pom.xml。3. ScribeSpanConsumer解码、过滤与异步入存储ScribeSpanConsumer.java 实现了生成的 Thrift 异步接口Scribe.AsyncIface是实际处理入口。Log(ListLogEntry messages, AsyncMethodCallbackResultCode resultHandler)的处理顺序为metrics.incrementMessages()记录收到一条消息遍历messages跳过 category 与配置不匹配的条目对匹配条目执行 Base64 MIME 解码与SpanBytesDecoder.THRIFT.decodeOne累加解码后二进制字节数到byteCount若解码过程抛出RuntimeException如非法 Base64则metrics.incrementMessagesDropped()并调用resultHandler.onError(e)最终在finally中metrics.incrementBytes(byteCount)解码成功的 Span 列表通过collector.accept(spans, callback, CommonPools.blockingTaskExecutor())异步提交因为存储组件可能不是异步的所以放在阻塞执行器上执行存储成功回调resultHandler.onComplete(ResultCode.OK)失败则回调onError。ResultCode与Scribe服务接口均为 Thrift 编译器0.12.0生成的代码见 generated 目录下的 Scribe.java 与 ResultCode.java。其中Scribe.AsyncIface.Log的签名即Log(ListLogEntry messages, AsyncMethodCallbackResultCode resultHandler)。4. 指标语义metrics在Builder.metrics()中被forTransport(scribe)作用域化因此与 Kafka、RabbitMQ 等其他采集器的指标相互独立。从 ScribeSpanConsumerTest.java 的断言中可以归纳出各指标的确切语义指标含义messages收到进入Log的 Scribe 消息条数messagesDropped解码/反序列化失败被丢弃的消息数bytes成功解码的 Span 二进制字节总数spans成功入队提交给存储的 Span 数spansDropped因采样或存储异常被丢弃的 Span 数测试覆盖了四条关键路径正常消费entriesWithSpansAreConsumed、category 不匹配被跳过entriesWithoutSpansAreSkipped、非法 Base64 被丢弃malformedDataIsDropped、存储侧异常不向客户端传播consumerExceptionBeforeCallbackDoesntSetFutureException以及 Finagle 客户端真实编码数据的兼容解码decodesSpanGeneratedByFinagle其 message 为带换行的 MIME Base64 文本块验证了 MIME 解码的必要性。端到端链路总览结合以上源码一条 Span 从 Scribe 客户端到存储的完整路径为Finagle/Scribe 客户端 → 写 TBinaryProtocol 大端序 Span 二进制 → Base64 编码 → scribe.log(categoryzipkin, messageencoded) → Netty LengthFieldBasedFrameDecoder解析 4 字节长度前缀帧 → ScribeInboundHandlerByteBuf → Armeria THttpService 请求 → ScribeSpanConsumer.Logcategory 过滤 → MIME Base64 解码 → THRIFT 反序列化 → Collector采样/指标 → SpanConsumer.accept 异步提交 → 存储后端Elasticsearch、Cassandra、MySQL 等由 zipkin-storage 提供其中Collector为 zipkin2 通用采集管道位于 zipkin-collector/core/src/main/java/zipkin2/collector/Collector.java负责采样、指标与存储调度ScribeCollector只负责协议接入这一层。测试与验证方式如果你要验证或调试本模块可以直接运行该模块的测试# 在仓库根目录执行模块内测试 ./mvnw -pl zipkin-collector/scribe test可参考的测试文件ScribeCollectorTest.java生命周期、匿名端口、端口冲突、toStringScribeSpanConsumerTest.java解码、过滤、异常路径与指标断言ITScribeCollector.java使用真实 Thrift 客户端TFramedTransportTBinaryProtocolScribe.Client向NettyScribeServer发送Log请求的集成测试。ITScribeCollector中的客户端构造方式同时也是自定义接入方如何向 Zipkin ScribeCollector 发送数据的最直接参考TTransport transport new TFramedTransport(new TSocket(localhost, server.port())); TProtocol protocol new TBinaryProtocol(transport, false, false); Scribe.Iface client new Scribe.Client(protocol); // 构造 LogEntry 列表后 ResultCode code client.Log(entries); assertThat(code).isEqualTo(ResultCode.OK);小结与适用前提ScribeCollector是一个面向历史 Finagle/Scribe 生态的兼容性采集器其价值在于让旧的 Scribe 上报链路无需改造即可汇入 Zipkin 的 zipkin2 存储与查询体系。接入时务必记住三点硬性约定编码Span 必须为TBinaryProtocol大端序二进制再经 Base64Basic 或 MIME 均可编码后放入LogEntry.message分类默认消费category zipkin可在Builder.category(...)或zipkin.collector.scribe.category中调整不匹配的条目会被静默跳过端口默认监听9410启用后即占用该 TCP 端口在 zipkin-server 中需显式设置zipkin.collector.scribe.enabledtrue才会加载。对于新接入方更推荐使用 Zipkin 原生支持的 HTTP/Kafka 等采集通道但如果你需要兼容存量 Scribe 客户端本模块开箱即用的 9410 端口、标准 Thrift 帧协议与完善的异常处理足以让集成过程平滑可控。赞分享可观测性后端微服务【免费下载链接】zipkinZipkin is a distributed tracing system项目地址https://gitcode.com/gh_mirrors/zip/zipkin点击查看免费下载相关推荐Zipkin Kafka Collector 深度指南基于 Kafka 0.10 的 Span 采集器配置、编码与源码剖析Zipkin Kafka Collector 深度指南基于 Kafka 0.10 的 Span 采集器配置、编码与源码剖析 导读 本文围绕 Zipkin 开可观测性后端微服务SkyWalking OTLP Trace 接入指南基于 OpenTelemetry 协议采集与 Zipkin 转换实战SkyWalking OTLP Trace 接入指南基于 OpenTelemetry 协议采集与 Zipkin 转换实战 SkyWalking 的 OAP 服可观测性APM链路追踪指标监控日志分析微服务Transform to Open Science (TOPS)开启21世纪科学革命的完整指南Transform to Open Science TOPS 开启21世纪科学革命的完整指南 Transform to Open Science TOPS 是数据库时序数据库物联网大数据实时分析云原生上一篇大语言模型对齐实战PPO与DPO算法原理及实现代码详解下一篇GoAccess终极命令行配置指南tmux与screen工作流优化创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表