.NET Core如何消费OTEL Agent推送至Kafka的Protobuf序列化日志
方案说明
1. OpenTelemetry Kafka Exporter JSON序列化配置
OTEL Kafka Exporter 原生支持JSON序列化,你只需要在kafka exporter配置中新增encoding参数即可,可选值为otlp_proto(默认Protobuf格式)、otlp_json(JSON格式)。
修改后的OTEL配置参考:
exporters: kafka: brokers: - "kafka:9093" protocol_version: 2.6.2 topic: logs encoding: otlp_json # 新增此行配置切换为JSON序列化 service: pipelines: logs: receivers: [filelog] exporters: [logging, kafka] # 记得把kafka加入日志pipeline的导出列表
配置生效后,Kafka中存储的就是结构化JSON日志,不管是.NET消费还是后续对接Kafka Connect同步Elasticsearch都不需要额外做序列化适配,是成本最低的方案。
2. 若需保留Protobuf序列化,.NET端消费方案
如果你需要保持Protobuf格式降低存储、传输成本,可以按以下步骤实现反序列化:
- 安装OTEL Protobuf定义的官方NuGet包,无需自己编写/编译Proto文件:
Install-Package OpenTelemetry.Proto.Logs - 基于Confluent.Kafka实现Protobuf日志反序列化,核心逻辑参考:
using Confluent.Kafka; using OpenTelemetry.Proto.Logs.V1; var config = new ConsumerConfig { GroupId = "otel-log-consumer-group", BootstrapServers = "kafka:9093", AutoOffsetReset = AutoOffsetReset.Earliest }; using var consumer = new ConsumerBuilder<Ignore, byte[]>(config).Build(); consumer.Subscribe("logs"); while (true) { var consumeResult = consumer.Consume(); // 反序列化Protobuf格式的OTEL日志 var logsData = LogsData.Parser.ParseFrom(consumeResult.Message.Value); // 处理日志逻辑,遍历日志资源、作用域、日志条目 foreach (var resourceLog in logsData.ResourceLogs) { foreach (var scopeLog in resourceLog.ScopeLogs) { foreach (var logRecord in scopeLog.LogRecords) { Console.WriteLine($"日志内容:{logRecord.Body.StringValue},时间:{logRecord.TimeUnixNano}"); } } } }
内容的提问来源于stack exchange,提问作者Ryan.Bartsch
相关产品推荐
相关产品推荐

