
SkyWalking 文件日志采集指南基于 Filebeat、Fluentd 与 Fluent-bit 的接入方案【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sk/skywalking应用日志是排查线上问题最重要的数据来源之一通常持久化在本地或网络文件系统中。SkyWalking OAP 本身不直接读取文件而是通过 Filebeat、Fluentd、Fluent-bit 等业界主流的日志采集器把文件日志转换成 SkyWalking 的日志数据协议再经 Kafka 或 HTTP 送入 OAP 分析、存储与检索。本文基于官方文档 Collecting File Log结合仓库中的 OAP 配置与端到端测试用例完整讲解三种采集器的接入步骤、日志格式化要求与协议细节帮助你在实际部署中把已有文件日志无缝接入 SkyWalking。采集整体架构SkyWalking 采集文件日志的链路分三段日志采集器Filebeat / Fluentd / Fluent-bit负责 tail 应用日志文件、解析出结构化字段并重写为 SkyWalking 的日志数据模型。传输通道KafkaFilebeat、Fluentd 使用或 OAP 的 HTTP REST 端口Fluent-bit 使用。OAP 后端通过kafka-fetcher消费 Kafka 中的日志或通过receiver-sharing-server的 HTTP 接口直接接收随后交给 log-analyzer 完成采样、分析、告警与存储。SkyWalking 定义了两类可直接被 OAP 消费的日志上报格式对应两种传输方式Kafka JSONnative-json通过 Kafka 上报参考 日志数据协议 - Native Kafka ProtocolHTTP JSON 数组通过 HTTP API 上报端点固定为http://oap-address:12800/v3/logs参考 日志数据协议 - HTTP API。一条完整的 JSON 日志记录包含以下核心字段以 Kafka/HTTP 两种格式一致字段必填说明timestamp可选日志产生时间毫秒时间戳缺省时 OAP 使用接收时间service必填逻辑服务名用于在拓扑图中形成独立节点serviceInstance可选服务实例名如 Pod 名或进程标识layer可选服务与实例的层如GENERAL缺省时 OAP 置为generaltraceContext可选追踪上下文traceId、traceSegmentId、spanId用于日志与 trace 关联tags可选可检索标签如level、logger等OAP 据此提供搜索/分析能力body必填日志正文支持text、json、yaml三种形态以 Kafka native-json 为例的完整记录结构完整定义见日志数据协议{ timestamp: 1618161813371, service: Your_ApplicationName, serviceInstance: 3a5b8da5a5ba40c0b192e91b5c80f1a8192.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} } }HTTP 方式则把上述记录放入 JSON 数组后 POST 到/v3/logs即可。理解这条记录结构是配置下文三种采集器的前提——所有采集器的工作本质都是把原始文本行转换成这个结构。准备开启 OAP 侧日志接收能力方式一通过 Kafka 接收Filebeat / FluentdKafka 采集链路依赖 OAP 的kafka-fetcher模块它默认处于关闭状态。在 application.yml 中按如下方式开启kafka-fetcher: selector: ${SW_KAFKA_FETCHER:default} default: bootstrapServers: ${SW_KAFKA_FETCHER_SERVERS:localhost:9092} namespace: ${SW_NAMESPACE:} partitions: ${SW_KAFKA_FETCHER_PARTITIONS:3} replicationFactor: ${SW_KAFKA_FETCHER_PARTITIONS_FACTOR:2} enableNativeProtoLog: ${SW_KAFKA_FETCHER_ENABLE_NATIVE_PROTO_LOG:true} enableNativeJsonLog: ${SW_KAFKA_FETCHER_ENABLE_NATIVE_JSON_LOG:true} consumers: ${SW_KAFKA_FETCHER_CONSUMERS:1}其中enableNativeJsonLog必须为true这是消费 native-json 格式文件日志的开关。从 KafkaFetcherConfig 源码可以看到enableNativeProtoLog与enableNativeJsonLog两个开关默认值均为true分别控制 proto 与 json 两类日志 topic 的消费。kafka-fetcher依赖的日志相关 topic 名为skywalking-logsnative-proto与skywalking-logs-jsonnative-json后者正是文件日志链路的目标 topictopic 名定义见 KafkaFetcherConfig。topic 不存在时 Kafka Fetcher 会按partitions与replicationFactor自动创建。更多 Kafka Fetcher 细节如命名空间隔离、MirrorMaker 2.0 配置参考 Kafka Fetcher 文档。方式二通过 HTTP 接收Fluent-bitHTTP 采集链路不需要 Kafka直接由 OAP 的 HTTP 服务接收日志。其监听地址为restHost:restPort有两个来源见 application.ymlreceiver-sharing-server模块启用时优先默认restHost: ${SW_RECEIVER_SHARING_REST_HOST:0.0.0.0}、restPort: ${SW_RECEIVER_SHARING_REST_PORT:0}core模块receiver-sharing-server未激活时兜底默认restHost: ${SW_CORE_REST_HOST:0.0.0.0}、restPort: ${SW_CORE_REST_PORT:12800}。即采集器把日志 POST 到restHost:restPort/v3/logs端点路径由共享 HTTP 服务统一提供。方案一Filebeat KafkaFilebeat 原生支持通过 Kafka 输出日志配置步骤为开启 OAP 的 kafka-fetcher 并确保enableNativeJsonLog: true见上文然后在 Filebeat 中完成「读取文件 → 解析 → 重组 → 输出 Kafka」。仓库在 test/e2e-v2/cases/kafka/log/filebeat.yml 提供了可直接参考的端到端测试配置核心结构如下http.enabled: true filebeat.inputs: - type: log enabled: true paths: - /tmp/skywalking-logs/*/e2e-service-provider.log processors: - dissect: tokenizer: [SW_CTX:%{SW_CTX}] [%{level}] %{logtime} [%{thread}] %{logger}:%{line} - %{body} field: message target_prefix: trim_values: all overwrite_keys: true - dissect: tokenizer: [%{service},%{serviceInstance},%{traceContext.traceId},%{traceContext.traceSegmentId},%{traceContext.spanId}] field: SW_CTX target_prefix: trim_values: all - rename: fields: - from: body to: body.text.text ignore_missing: false fail_on_error: true - drop_fields: fields: [ecs, host, input, agent, log, message, SW_CTX, level, thread, logger, line, logtime] output.kafka: hosts: [broker-a:9092, broker-b:9092] topic: skywalking-logs-json要点解析两段式 dissect 解析先用外层 tokenizer 切出SW_CTX上下文块与level、logtime、thread、logger、line、body等日志字段再用内层 tokenizer 从SW_CTX中切出service、serviceInstance与 trace 上下文三个 ID。这种嵌套写法对应 Java Agent 通过日志插件向日志文本注入的[SW_CTX:...]前缀traceContext各字段名与协议中的TraceContext一一对应。字段重组将body重命名到body.text.text从而生成协议要求的LogDataBody结构drop_fields清理中间字段避免把message、log等原始字段带入 Kafka。时间戳处理测试用例还通过 javascriptscript处理器把logtime如2021-04-11 17:23:33.371解析为毫秒时间戳写入timestamp字段生产环境可根据自身日志格式实现等价转换。输出目标output.kafka的topic必须为skywalking-logs-json与 OAP 端topicNameOfJsonLogs保持一致hosts指向实际 Kafka 集群。方案二Fluentd KafkaFluentd 同样通过 Kafka 上报前置条件与 Filebeat 一致开启 kafka-fetcher、enableNativeJsonLog: true。参考配置位于 test/e2e-v2/cases/kafka/log/fluentd.conf分段说明如下。1. 输入tail 文件并按正则解析source type tail path /tmp/skywalking-logs/*/e2e-service-provider.log read_from_head true tag skywalking_log refresh_interval 1 parse type regexp expression /^\[SW_CTX:\s*(?SW_CTX[^ ]*]*)\] \[(?level[^ ]*)\] (?logtime[^\]]*) \[(?thread[^ ]*)\] (?logger[^\]]*):(?line[^\]]*) - (?body[^\]]*)$/ /parse /sourcepath支持通配符匹配多个日志文件正则与 Filebeat 的 dissect tokenizer 对应负责把一行日志拆成命名捕获组。2. 过滤解析 SW_CTX 上下文filter skywalking_log type parser key_name SW_CTX reserve_data true remove_key_name_field true parse type regexp expression /^\[(?service[^\]]*),(?serviceInstance[^\]]*),(?traceId[^\]]*),(?traceSegmentId[^\]]*),(?spanId[^\]]*)\]$/ /parse /filter把SW_CTX中的服务、实例与 trace 三个 ID 提取为独立字段。3. 过滤重组为协议结构filter skywalking_log type record_transformer enable_ruby true record traceContext ${{traceId record[traceId], traceSegmentId record[traceSegmentId], spanId record[spanId]}} tags ${{data [{key level, value record[level]}, {key thread, value record[thread]}, {key logger, value record[logger]}, {key agent, value fluentd}]}} body ${{text {text record[body]}}} timestamp ${DateTime.strptime(record[logtime], %Y-%m-%d %H:%M:%S.%L).strftime(%Q)} /record remove_keys level,thread,logger,line,traceId,traceSegmentId,spanId,logtime /filter这里用 Ruby 表达式完成三件关键事把三个 ID 组装成traceContext对象把level、thread、logger等组装成tags.data标签数组并补充agent: fluentd标识采集器来源把正文包进body.text.text。timestamp用strptime把logtime解析成毫秒时间戳与协议字段对齐。4. 输出写入 Kafka JSON topicmatch skywalking_log type kafka2 brokers broker-a:9092,broker-b:9092 default_topic skywalking-logs-json format type json /format /matchdefault_topic固定为skywalking-logs-json序列化格式为 JSON与 OAP 的 native-json 解析器完全匹配。方案三Fluent-bit HTTP与前两者不同Fluent-bit不经过 Kafka而是把日志直接 POST 到 OAP 的 HTTP REST 端口即receiver-sharing-server的restHost:restPort若该模块未启用则落到core的restPort:12800见上文「方式二」。仓库在 test/e2e-v2/cases/log/fluent-bit/ 提供了完整示例共三份文件配合使用fluent-bit.conf主配置[SERVICE] Flush 5 Daemon Off Log_Level warn Parsers_File fluent-bit-parser.conf [INPUT] Name tail Path /tmp/skywalking-logs/*/e2e-service-provider.log Parser my-log-format [FILTER] Name lua Match * Script fluent-bit-script.lua Call rewrite_body [OUTPUT] Name stdout Match * Format json [OUTPUT] Name http Match * Host oap Port 12800 URI /v3/logs Format jsonParsers_File引入解析器定义文件INPUT中Parser my-log-format指定行解析器FILTER通过 Lua 脚本把解析出的字段重写为 SkyWalking 协议结构两个OUTPUTstdout用于调试查看最终 JSONhttp是正式出口——Host指向 OAP 地址、Port为restPort、URI固定为/v3/logs格式为 JSON。fluent-bit-parser.conf正则解析器[PARSER] Name my-log-format Format regex Regex ^\[SW_CTX: ?\[(?service[^,]),(?serviceInstance[^,]),(?traceId[^,]),(?traceSegmentId[^,]),(?spanId[^\]])\]\] \[(?level.?)\] (?logtime[^\]]*) \[(?thread[^ ]*)\] (?logger[^\]]*):(?line[^\]]*) - (?body[^\]]*)$ Time_Key time Time_Format %d/%b/%Y:%H:%M:%S %z该正则一次性把整行日志含SW_CTX嵌套括号拆成service、serviceInstance、traceId、traceSegmentId、spanId、level、logtime、thread、logger、line、body等字段。fluent-bit-script.lua字段重组脚本function rewrite_body(tag, timestamp, record) record[body] {text{textfluentbit .. record[body]}} record[tags] {data{{keylevel, valueINFO}}} record[traceContext] {traceIdrecord[traceId],traceSegmentIdrecord[traceSegmentId],spanIdrecord[spanId]} return 1, timestamp, record endLua 脚本与前面 Filebeat/Fluentd 的 processor/filter 作用相同把正文包成body.text.text、把level等组装成tags.data、把三个 ID 组装成traceContext最终由 HTTP output 以 JSON 数组形式上报。测试用例中该脚本把tags固定为INFO级别、正文加了fluentbit前缀实际使用时请替换为从解析字段中读取的真实值。三种方案对比与选型建议维度Filebeat KafkaFluentd KafkaFluent-bit HTTPOAP 前置条件开启kafka-fetcherenableNativeJsonLog: true同左无需 Kafkareceiver-sharing-server或core的 REST 端口可达传输协议Kafkatopicskywalking-logs-jsonKafkatopicskywalking-logs-jsonHTTP POST/v3/logs解析手段dissectscript(javascript)regexprecord_transformer(ruby)regexparser luafilter是否支持 trace 关联支持解析SW_CTX注入traceContext支持支持典型场景已有 Elastic 技术栈、希望复用电采体系已有 Fluentd 采集生态含插件扩展轻量级部署、希望最小化组件依赖三条链路对 OAP 是等价的最终都以同一套日志数据模型进入 log-analyzer。选择依据主要是你现有的采集器生态与是否愿意引入 Kafka。若集群中已部署 KafkaFilebeat/Fluentd 方案可复用 topic 与消费体系便于统一管理若追求组件最少、链路最短Fluent-bit 直连 HTTP 更直接也适合容器环境下以 sidecar 方式采集。验证与排障要点调试输出Fluent-bit 的stdoutoutput 可在正式上链路前打印最终 JSON对照上文协议示例逐字段核对Filebeat/Fluentd 也可先用 stdout/type copy输出验证解析结果。Topic 校验Kafka 方案下确认消息确实进入skywalking-logs-jsontopic可用kafka-console-consumer抽查且 OAP 日志中无消费异常。enableNativeJsonLog检查若 Kafka 中消息堆积但 OAP 无日志入库优先确认application.yml中kafka-fetcher.default.enableNativeJsonLog为true。HTTP 端点校验Fluent-bit 方案确认 OAP 的restPort可被采集器访问URI 为/v3/logs日志记录需以 JSON 数组形式提交参考 HTTP API 示例。时间戳一致性协议中timestamp为毫秒时间戳三份示例配置都通过脚本/表达式把日志时间字符串转为毫秒格式不一致会导致时间错乱或入库异常。安全与限流生产环境建议结合 OAP 的 token 认证参考 backend-token-auth保护 HTTP 日志接入端口并评估日志量对 OAP 与存储如 Elasticsearch的写入压力。完成上述配置后文件日志即可在 SkyWalking UI 的 Log 页面按服务、实例、level、logger等标签检索并可与对应 trace 直接关联跳转形成「trace log」一体化的排障闭环。【免费下载链接】skywalkingAPM, Application Performance Monitoring System项目地址: https://gitcode.com/gh_mirrors/sk/skywalking创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考