SkyWalking 日志数据上报协议全解:gRPC / Kafka / HTTP 三种方式实战指南
2026/9/20 14:37:18 网站建设 项目流程

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 上报、高吞吐流式推送
Kafkanative-json(JSON)skywalking-logs-jsonTopic日志先聚合到 Kafka,再由 OAP 异步消费
HTTP APIJSON 数组POST http://<oap-address>:12800/v3/logs简单集成、无 Agent 场景、脚本上报

尽管传输层各不相同,三条通道共享同一份LogData数据模型:OAP 侧统一将它们转换为内部日志元数据并交给 LAL(Log Analysis Language)分析器处理。因此,理解LogData的字段语义是掌握整个协议的关键。

说明:Kafka 通道同时也支持 protobuf 二进制格式(默认 Topicskywalking-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 注入日志文本后,上报时携带traceIdtraceSegmentIdspanId,可实现日志与链路追踪的关联查询
  • tags(可选):键值对标签,OAP 基于这些标签提供搜索/分析能力(如levellogger等)。
  • layer(可选,自 9.0.0 起):服务与实例的层(layer)。缺省时 OAP 自动设置为layer = ID: 2, NAME: general。常见的层还有GENERALMESHVIRTUAL_DATABASEFAAS等,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 规则中正则或拆分机制提取有意义的信息(如jsonregex解析器),参见 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: INFOlogger: com.example.MyLogger

2.5 数据到达 OAP 后的流转:源码级佐证

三条通道的日志最终都会汇入同一个分析入口。以 gRPC 通道为例,LogReportServiceGrpcHandler.java 实现了collect双向流:

  • 每次onNext(LogData)时,先执行setServiceName(builder)若当前流已缓存了 serviceName,就用缓存值覆盖消息中的空 service 字段——这正是文档中"service 可共享前一条值"的落地实现;
  • 随后调用logAnalyzerService.doAnalysis(LogMetadataUtils.fromLogData(builder), builder)进入 LAL 分析;
  • 流结束后回写空的CommandsonCompleted()

LogMetadataUtils.java 负责把LogData中的serviceserviceInstanceendpointlayertimestamptraceContext提取为内部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)模式上报日志,理由有二:

  1. 字段继承:同一流中未显式设置的service/serviceInstance/endpoint可继承前一条的值,减少重复字段传输;
  2. 网络成本:批量上报同一服务的日志可显著降低网络开销。

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 中topicNameOfLogstopicNameOfJsonLogs的默认值,可通过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
  • bodytext.text承载实际日志内容,也可替换为jsonyaml结构。

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 通道不做流式字段继承,每条日志都应携带完整的关键字段(至少servicebody)。

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):选择 Kafkanative-json通道,将 JSON 序列化的LogData写入skywalking-logs-json,OAP 异步消费,天然解耦;
  • 轻量集成、脚本或临时验证:选择 HTTP API,一条curl即可完成。

6.2 数据到达后的 LAL 分析

无论走哪条通道,日志最终都会进入ILogAnalyzerService.doAnalysis,由 LAL(Log Analysis Language)规则决定如何解析、清洗、提取指标与存储。协议中LogDataBody.typetagslayer等字段正是为 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 →LogMetadataUtilsILogAnalyzerService.doAnalysis),开发者既可以快速实现自定义上报端,也能在排查接入问题时准确定位链路节点。

【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sk/skywalking

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询