SkyWalking 日志数据上报协议全解:gRPC / Kafka / HTTP 三种方式实战指南
【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sk/skywalking
导读:本文是 Apache SkyWalking OAP 日志数据接入的权威协议指南。围绕 docs/en/api/log-data-protocol.md 展开,完整讲解
LogData核心数据结构,以及gRPC(native-proto)、Kafka(native-json)、HTTP(JSON)三种上报通道的格式规范、字段语义与接入要点。读完本文,你将能够自行编写上报客户端或改造日志采集链路,将任意日志源(应用日志、文件抓取、Satellite 采集)以标准协议接入 SkyWalking OAP,并理解 OAP 侧从接收入口到 LAL 分析的完整调用链。
一、协议总览:三种通道、一份数据模型
SkyWalking 为日志数据提供了三条上报通道,分别面向不同场景:
| 通道 | 数据格式 | 入口 | 典型场景 |
|---|---|---|---|
| gRPC 流式 | native-proto(Protobuf) | LogReportService.collect双向流 RPC | 语言 Agent 上报、高吞吐流式推送 |
| Kafka | native-json(JSON) | skywalking-logs-jsonTopic | 日志先聚合到 Kafka,再由 OAP 异步消费 |
| HTTP API | JSON 数组 | POST http://<oap-address>:12800/v3/logs | 简单集成、无 Agent 场景、脚本上报 |
尽管传输层各不相同,三条通道共享同一份LogData数据模型:OAP 侧统一将它们转换为内部日志元数据并交给 LAL(Log Analysis Language)分析器处理。因此,理解LogData的字段语义是掌握整个协议的关键。
说明:Kafka 通道同时也支持 protobuf 二进制格式(默认 Topic
skywalking-logs),native-json是其中一种数据格式,下文详细展开。
二、核心数据模型:LogData 字段深度解析
LogData是日志上报协议的统一消息体,其 Protobuf 定义(package skywalking.v3)完整如下:
syntax = "proto3"; package skywalking.v3; option java_multiple_files = true; option java_package = "org.apache.skywalking.apm.network.logging.v3"; option csharp_namespace = "SkyWalking.NetworkProtocol.V3"; option go_package = "skywalking.apache.org/repo/goapi/collect/logging/v3"; import "common/Common.proto"; import "common/Command.proto"; // Report collected logs into the OAP backend service LogReportService { // Recommend to report log data in a stream mode. // The service/instance/endpoint of the log could share the previous value if they are not set. // Reporting the logs of same service in the batch mode could reduce the network cost. rpc collect (stream LogData) returns (Commands) { } } // Log data is collected through file scratcher of agent. // Natively, Satellite provides various ways to collect logs. message LogData { // [Optional] The timestamp of the log, in millisecond. // If not set, OAP server would use the received timestamp as log's timestamp, or relies on the OAP server analyzer. int64 timestamp = 1; // [Required] **Service**. Represents a set/group of workloads which provide the same behaviours for incoming requests. // // The logic name represents the service. This would show as a separate node in the topology. // The metrics analyzed from the spans, would be aggregated for this entity as the service level. // // If this is not the first element of the streaming, use the previous not-null name as the service name. string service = 2; // [Optional] **Service Instance**. Each individual workload in the Service group is known as an instance. Like `pods` in Kubernetes, it // doesn't need to be a single OS process, however, if you are using instrument agents, an instance is actually a real OS process. // // The logic name represents the service instance. This would show as a separate node in the instance relationship. // The metrics analyzed from the spans, would be aggregated for this entity as the service instance level. string serviceInstance = 3; // [Optional] **Endpoint**. A path in a service for incoming requests, such as an HTTP URI path or a gRPC service class + method signature. // // The logic name represents the endpoint, which logs belong. string endpoint = 4; // [Required] The content of the log. LogDataBody body = 5; // [Optional] Logs with trace context TraceContext traceContext = 6; // [Optional] The available tags. OAP server could provide search/analysis capabilities based on these. LogTags tags = 7; // [Optional] Since 9.0.0 // The layer of the service and servce instance. If absent, the OAP would set `layer`=`ID: 2, NAME: general` string layer = 8; }2.1 核心字段语义与取值要点
timestamp(可选,毫秒):日志产生时间戳。不设置时,OAP 会使用服务端接收时间作为日志时间,或交由 OAP 的 analyzer 决定,因此高精度日志场景建议显式填充。service(必填):服务逻辑名,代表一组行为相同的工作负载。它会作为拓扑图中的一个独立节点出现,日志指标也以该实体为粒度做服务级聚合。流式上报时若后续消息未填,沿用同一流中前一条非空值。serviceInstance(可选):服务实例,Service 组内的每个独立工作负载,类似 Kubernetes 中的pods。使用探针 Agent 时一个实例通常对应一个真实 OS 进程。它会在实例关系图中作为独立节点展示,指标按实例级聚合。endpoint(可选):端点,即服务内接收请求的路径,例如 HTTP URI 路径或 gRPC 的"类 + 方法签名"。表示日志归属的端点。body(必填):日志正文,通过LogDataBody携带,支持文本 / JSON / YAML 三种格式(详见第三节)。traceContext(可选):Trace 上下文。当探针把 Trace ID 注入日志文本后,上报时携带traceId、traceSegmentId、spanId,可实现日志与链路追踪的关联查询。tags(可选):键值对标签,OAP 基于这些标签提供搜索/分析能力(如level、logger等)。layer(可选,自 9.0.0 起):服务与实例的层(layer)。缺省时 OAP 自动设置为layer = ID: 2, NAME: general。常见的层还有GENERAL、MESH、VIRTUAL_DATABASE、FAAS等,OAP 据此匹配对应的 LAL 规则与存储模板。
2.2 正文载体:LogDataBody 的三种格式
LogDataBody使用type字段(字符串)匹配 OAP 侧 analyzer,content为可扩展的 oneof 字段:
// The content of the log data message LogDataBody { // A type to match analyzer(s) at the OAP server. // The data could be analyzed at the client side, but could be partial string type = 1; // Content with extendable format. oneof content { TextLog text = 2; JSONLog json = 3; YAMLLog yaml = 4; } } // Literal text log, typically requires regex or split mechanism to filter meaningful info. message TextLog { string text = 1; } // JSON formatted log. The json field represents the string that could be formatted as a JSON object. message JSONLog { string json = 1; } // YAML formatted log. The yaml field represents the string that could be formatted as a YAML map. message YAMLLog { string yaml = 1; }TextLog:纯文本日志。通常需要 OAP 侧 LAL 规则中正则或拆分机制提取有意义的信息(如json、regex解析器),参见 lal.md。JSONLog:JSON 格式日志,json字段是可解析为 JSON 对象的字符串,便于 LAL 直接按 JSON 路径取字段。YAMLLog:YAML 格式日志,yaml字段是能解析为 YAML 映射的字符串。
2.3 关联 Trace:TraceContext
// Logs with trace context, represent agent system has injects context(IDs) into log text. message TraceContext { // [Optional] A string id represents the whole trace. string traceId = 1; // [Optional] A unique id represents this segment. Other segments could use this id to reference as a child segment. string traceSegmentId = 2; // [Optional] The number id of the span. Should be unique in the whole segment. // Starting at 0. int32 spanId = 3; }2.4 标签容器:LogTags
message LogTags { // String key, String value pair. repeated KeyStringValuePair data = 1; }LogTags复用common/Common.proto中的KeyStringValuePair,即key/value字符串对列表。这些标签直接决定日志可搜索与可分析的维度,常见用法如level: INFO、logger: com.example.MyLogger。
2.5 数据到达 OAP 后的流转:源码级佐证
三条通道的日志最终都会汇入同一个分析入口。以 gRPC 通道为例,LogReportServiceGrpcHandler.java 实现了collect双向流:
- 每次
onNext(LogData)时,先执行setServiceName(builder):若当前流已缓存了 serviceName,就用缓存值覆盖消息中的空 service 字段——这正是文档中"service 可共享前一条值"的落地实现; - 随后调用
logAnalyzerService.doAnalysis(LogMetadataUtils.fromLogData(builder), builder)进入 LAL 分析; - 流结束后回写空的
Commands并onCompleted()。
LogMetadataUtils.java 负责把LogData中的service、serviceInstance、endpoint、layer、timestamp、traceContext提取为内部LogMetadata,作为 LAL 脚本的输入上下文。HTTP 通道的 LogReportServiceHTTPHandler.java 与之完全相同——三条通道殊途同归,都汇聚到ILogAnalyzerService.doAnalysis。
此外,gRPC 与 HTTP 两个 Handler 都注册了遥测指标:log_in_latency(日志处理延迟直方图)与log_analysis_error_count(分析错误计数),并以protocol(grpc/http/kafka)与data_format(protobuf/json)作为标签区分通道,便于在 backend-telemetry.md 中观测各通道健康度。
三、Native Proto Protocol:gRPC 流式上报
3.1 服务定义与推荐用法
gRPC 通道由LogReportService服务提供(完整定义见上文Logging.proto,即 apm-protocol 相关协议目录 所生成代码的对应源):
service LogReportService { rpc collect (stream LogData) returns (Commands) { } }官方推荐以流式(stream)模式上报日志,理由有二:
- 字段继承:同一流中未显式设置的
service/serviceInstance/endpoint可继承前一条的值,减少重复字段传输; - 网络成本:批量上报同一服务的日志可显著降低网络开销。
3.2 流式请求的字段继承机制
结合 LogReportServiceGrpcHandler.java 的实现细节,字段继承规则为:
- 流的首个消息必须携带
service; - 后续消息若
service为空,Handler 会用serviceName缓存覆盖; - 该缓存仅对当前这条 gRPC 流生效,流结束后即失效。
因此业务上建议按服务维度建立独立的流,既能利用继承减少 payload,也能保证故障隔离。
四、Native Kafka Protocol:经 Kafka 上报 native-json
4.1 通道说明与 Topic 约定
Kafka 通道用于上报native-json格式日志。OAP 的kafka-fetcher负责消费 Kafka 中的日志数据,相关 Topic 约定见 kafka-fetcher.md:
skywalking-logs:native proto(protobuf)格式日志;skywalking-logs-json:native json 格式日志(本文档对应的通道)。
上述 Topic 名称同时作为 KafkaFetcherConfig.java 中topicNameOfLogs、topicNameOfJsonLogs的默认值,可通过kafka-fetcher配置覆盖。
4.2 Kafka JSON 日志记录示例(完整)
Kafka 消息体(value)为一条 JSON 序列化后的LogData:
{ "timestamp":1618161813371, "service":"Your_ApplicationName", "serviceInstance":"3a5b8da5a5ba40c0b192e91b5c80f1a8@192.168.1.8", "layer":"GENERAL", "traceContext":{ "traceId":"ddd92f52207c468e9cd03ddd107cd530.69.16181331190470001", "spanId":"0", "traceSegmentId":"ddd92f52207c468e9cd03ddd107cd530.69.16181331190470000" }, "tags":{ "data":[ { "key":"level", "value":"INFO" }, { "key":"logger", "value":"com.example.MyLogger" } ] }, "body":{ "text":{ "text":"log message" } } }要点提示:
- JSON 字段名与 Protobuf 字段名一一对应,
timestamp为毫秒整数,spanId在 JSON 中表现为字符串(示例为"0"); layer显式指定为GENERAL;若省略,OAP 会按文档约定补为general;body中text.text承载实际日志内容,也可替换为json或yaml结构。
4.3 OAP 侧消费实现
JsonLogHandler.java 负责消费skywalking-logs-jsonTopic:它将 Kafka 消息 value 按 UTF-8 解码后,通过ProtoBufJsonUtils.fromJSON反序列化为LogData.Builder,再走与 gRPC 相同的logAnalyzerService.doAnalysis分析链路(继承自 LogHandler.java)。LogHandler同时消费 protobuf 格式(默认 Topicskywalking-logs),两种格式共用同一分析入口,只是data_format遥测标签不同。
五、HTTP API:POST /v3/logs
5.1 端点与请求方式
HTTP 通道端点固定为:
http://<oap-address>:12800/v3/logs- 方法:
POST; - 请求体:JSON 数组,每个元素是一条
LogData的 JSON 表示; - 响应:
Commands(空命令列表)或异常。
该端点由 LogReportServiceHTTPHandler.java 中的@Post("/v3/logs")注解注册到共享 HTTP 服务端口 12800(SharingServerModule),注册逻辑见 LogModuleProvider.java。Handler 逐条将日志转为LogData.Builder并调用doAnalysis,处理完成后返回空Commands。
5.2 完整请求示例
[ { "timestamp": 1618161813371, "service": "Your_ApplicationName", "serviceInstance": "3a5b8da5a5ba40c0b192e91b5c80f1a8@192.168.1.8", "layer":"GENERAL", "traceContext": { "traceId": "ddd92f52207c468e9cd03ddd107cd530.69.16181331190470001", "spanId": "0", "traceSegmentId": "ddd92f52207c468e9cd03ddd107cd530.69.16181331190470000" }, "tags": { "data": [ { "key": "level", "value": "INFO" }, { "key": "logger", "value": "com.example.MyLogger" } ] }, "body": { "text": { "text": "log message" } } } ]字段语义与 Kafka JSON 记录完全一致,唯一区别是 HTTP 一次可提交多条日志的数组(上例为单条,实际可批量)。需要注意:HTTP 通道不做流式字段继承,每条日志都应携带完整的关键字段(至少service与body)。
5.3 快速验证
可用curl做连通性验证:
curl -X POST http://<oap-address>:12800/v3/logs \ -H 'Content-Type: application/json' \ -d '[ { "timestamp": 1618161813371, "service": "Your_ApplicationName", "serviceInstance": "3a5b8da5a5ba40c0b192e91b5c80f1a8@192.168.1.8", "layer": "GENERAL", "body": { "text": { "text": "hello skywalking log" } } } ]'返回空Commands即表示 OAP 已接收;随后可在 UI 的日志查询页面(Log Query)中检索该service的日志,验证链路是否打通。
六、接入决策与后续处理链路
6.1 三条通道如何选择
- 有探针/自定义 Agent:优先 gRPC 流式上报,利用字段继承与批处理降低网络成本;
- 已有 Kafka 日志汇聚(如 Filebeat → Kafka):选择 Kafka
native-json通道,将 JSON 序列化的LogData写入skywalking-logs-json,OAP 异步消费,天然解耦; - 轻量集成、脚本或临时验证:选择 HTTP API,一条
curl即可完成。
6.2 数据到达后的 LAL 分析
无论走哪条通道,日志最终都会进入ILogAnalyzerService.doAnalysis,由 LAL(Log Analysis Language)规则决定如何解析、清洗、提取指标与存储。协议中LogDataBody.type、tags、layer等字段正是为 LAL 脚本提供上下文的关键输入。更多 LAL 语法与解析器(regex / json / grok 等)参见 lal.md 与示例配置 lal.yaml;日志分析功能的启用与配置参见 log-analyzer.md。
6.3 相关文档导航
- 日志采集与分析整体说明:log-analyzer.md
- LAL 规则语言:lal.md
- Kafka 通道配置(Topic、消费者等):kafka-fetcher.md
- 日志配置示例(告警、模板):alarm-settings.yml、log-mal.yaml
- OAP 遥测指标说明:backend-telemetry.md
七、结语
SkyWalking 的日志数据协议以LogData为统一模型,通过 gRPC、Kafka、HTTP 三条通道覆盖了从 Agent 直连到消息队列中转、再到轻量 API 集成的全部典型场景。协议设计上的几个关键点值得在实际接入时牢记:流式上报的字段继承机制、layer缺省为general的语义、tags决定日志可搜索维度、traceContext打通日志与 Trace 关联。结合本文提供的源码调用链(三个 Handler →LogMetadataUtils→ILogAnalyzerService.doAnalysis),开发者既可以快速实现自定义上报端,也能在排查接入问题时准确定位链路节点。
【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sk/skywalking
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考