Go日志采集:ELK与Loki集成
Go日志采集:ELK与Loki集成摘要: 本篇讲解Go语言日志采集实战配置zap输出JSON日志供Filebeat采集到Elasticsearch搭建LokiPromtail日志体系设计带trace_id的日志结构实现全流程追踪用LogQL和Kibana查询日志分享日志量过大导致Elasticsearch内存溢出的踩坑经验对比ELK和Loki两种日志平台的架构与成本。开篇故事去年我们微服务集群从8个服务扩展到30个日志分散在各台机器上。每次排查线上问题要SSH登6-7台机器grep日志一个跨服务的问题排查半小时起步。当时决定上ELK搭了Elasticsearch集群加Filebeat采集第二天日志就能在Kibana里全文检索了排查效率提升明显。但跑了两个月问题来了。30个服务每天产生50GB日志Elasticsearch内存占用超过32GB节点频繁Full GC。加机器也不行ES存原始日志太费空间。后来引入Loki只索引标签不索引日志正文同样的日志量Loki只占ES的十分之一存储。这篇把ELK和Loki两套体系的集成方法写清楚。一、zap输出JSON供Filebeat采集日志采集的第一步是让应用输出结构化日志。上一篇讲了zap的JSON编码这里配置一个完整的生产级日志器确保输出格式能被Filebeat正确解析。packagemainimport(contextcrypto/randencoding/hexnet/httpostimego.uber.org/zapgo.uber.org/zap/zapcore)// LogConfig 日志配置typeLogConfigstruct{Levelstring// 日志级别: debug/info/warn/errorFilePathstring// 日志文件路径MaxSizeint// 单文件最大MBMaxBackupsint// 保留旧文件数MaxAgeint// 保留天数Compressbool// 是否压缩旧日志}// NewProductionLogger 创建生产级Logger// 输出JSON格式带trace_id和service_namefuncNewProductionLogger(cfg LogConfig)(*zap.Logger,error){// 日志文件轮转配置// lumberjack实现日志文件按大小切割和保留// 这里用os.OpenFile简化实际生产用lumberjackfile,err:os.OpenFile(cfg.FilePath,os.O_CREATE|os.O_WRONLY|os.O_APPEND,0644,)iferr!nil{returnnil,err}// 自定义Encoder输出Filebeat能解析的JSONencoderConfig:zapcore.EncoderConfig{TimeKey:timestamp,// Filebeat默认识别的字段名LevelKey:level,MessageKey:message,// Filebeat默认识别的字段名CallerKey:caller,StacktraceKey:stack_trace,LineEnding:zapcore.DefaultLineEnding,// 时间用RFC3339Filebeat自动解析为时间类型EncodeTime:zapcore.RFC3339NanoTimeEncoder,// 级别用小写符合Elasticsearch习惯EncodeLevel:zapcore.LowercaseLevelEncoder,EncodeCaller:zapcore.ShortCallerEncoder,}// 解析日志级别varlevel zapcore.Leveliferr:level.UnmarshalText([]byte(cfg.Level));err!nil{returnnil,err}core:zapcore.NewCore(zapcore.NewJSONEncoder(encoderConfig),zapcore.AddSync(file),level,)// 添加公共字段: 服务名和实例ID// 每条日志都带这两个字段方便在ELK中过滤logger:zap.New(core,zap.AddCaller(),zap.AddCallerSkip(1),// With方法添加公共字段后续每条日志都会带上).With(zap.String(service,order-service),zap.String(instance,getHostname()),)returnlogger,nil}// getHostname 获取主机名funcgetHostname()string{host,err:os.Hostname()iferr!nil{returnunknown}returnhost}// TraceLogger 创建带trace_id的Logger// trace_id贯穿整个请求过程串联多个服务的日志funcTraceLogger(base*zap.Logger,traceID,spanIDstring)*zap.Logger{returnbase.With(zap.String(trace_id,traceID),zap.String(span_id,spanID),)}// LogMiddleware HTTP请求日志中间件// 每个请求生成trace_id注入到日志和响应头funcLogMiddleware(logger*zap.Logger,next http.Handler)http.Handler{returnhttp.HandlerFunc(func(w http.ResponseWriter,r*http.Request){// 生成或从请求头提取trace_idtraceID:r.Header.Get(X-Trace-Id)iftraceID{traceIDgenerateTraceID()}// 注入到响应头前端可以拿到w.Header().Set(X-Trace-Id,traceID)// 创建带trace_id的Loggerlog:TraceLogger(logger,traceID,generateSpanID())// 把Logger和trace_id存入contextctx:context.WithValue(r.Context(),loggerKey{},log)ctxcontext.WithValue(ctx,traceIDKey{},traceID)// 记录请求开始start:time.Now()log.Info(request started,zap.String(method,r.Method),zap.String(path,r.URL.Path),zap.String(remote_addr,r.RemoteAddr),)// 包装ResponseWriter捕获状态码wrapped:statusWriter{ResponseWriter:w,status:200}next.ServeHTTP(wrapped,r.WithContext(ctx))// 记录请求结束log.Info(request completed,zap.Int(status,wrapped.status),zap.Duration(duration,time.Since(start)),)})}// 从context获取Logger的工具typeloggerKeystruct{}typetraceIDKeystruct{}// FromContext 从context获取LoggerfuncFromContext(ctx context.Context)*zap.Logger{ifl,ok:ctx.Value(loggerKey{}).(*zap.Logger);ok{returnl}returnzap.NewNop()// 没有Logger返回空操作Logger}// statusWriter 包装ResponseWritertypestatusWriterstruct{http.ResponseWriter statusint}func(sw*statusWriter)WriteHeader(codeint){sw.statuscode sw.ResponseWriter.WriteHeader(code)}// generateTraceID 生成trace IDfuncgenerateTraceID()string{b:make([]byte,16)rand.Read(b)returnhex.EncodeToString(b)}funcgenerateSpanID()string{b:make([]byte,8)rand.Read(b)returnhex.EncodeToString(b)}日志中间件给每个请求分配trace_id后续同一次请求的所有日志都带这个ID。在ELK或Loki里搜索trace_id:abc123就能看到一次请求经过所有服务的完整日志。二、Filebeat采集与Loki Promtail配置日志文件准备好后需要采集器把日志送到Elasticsearch或Loki。Filebeat是ELK体系的采集器Promtail是Loki体系的采集器。Filebeat配置文件filebeat.yml核心部分:filebeat.inputs:-type:logpaths:[/var/log/go/*.log]# JSON解析器直接解析zap输出的JSON日志json.keys_under_root:true# 多行配置: 堆栈信息合并到上一行multiline.pattern:^\smultiline.match:afteroutput.elasticsearch:hosts:[es-node1:9200]# 按级别分索引indices:-index:go-logs-error-%{yyyy.MM.dd}when.contains:{level:error}-index:go-logs-info-%{yyyy.MM.dd}when.contains:{level:info}# 索引生命周期: 30天后自动删除setup.ilm.enabled:truesetup.ilm.rollover_alias:go-logsLoki的Promtail配置文件promtail.yml核心部分:# Promtail采集Go服务日志到Lokiclients:-url:http://loki:3100/loki/api/v1/pushscrape_configs:-job_name:go_servicestatic_configs:-targets:[localhost]labels:job:go-serviceenv:production__path__:/var/log/go/*.log# JSON解析pipelinepipeline_stages:# 解析JSON日志-json:expressions:level:levelservice:servicetrace_id:trace_id# 提取level和service作为标签-labels:level:service:# 提取时间戳-timestamp:source:timestampformat:RFC3339NanoLoki和ES的核心区别: ES索引全文Loki只索引标签。Loki的labels阶段把level和service提取为标签只有这些字段能快速过滤。日志正文不建索引查询时顺序扫描。所以Loki存储成本低但全文检索比ES慢。三、日志检索查询日志采到ELK和Loki后怎么查询。两种平台的查询语言不同。packagemainimport(contextencoding/jsonfmtnet/httpnet/urlstringstime)// LokiClient Loki日志查询客户端typeLokiClientstruct{baseURLstringhttpClient*http.Client}// NewLokiClient 创建Loki客户端funcNewLokiClient(baseURLstring)*LokiClient{returnLokiClient{baseURL:strings.TrimRight(baseURL,/),httpClient:http.Client{Timeout:30*time.Second,},}}// QueryRange 范围日志查询// LogQL语法: {serviceorder-service} | error | trace_idabc123func(c*LokiClient)QueryRange(ctx context.Context,querystring,start,end time.Time,limitint)([]LogEntry,error){params:url.Values{}params.Set(query,query)params.Set(start,fmt.Sprintf(%d,start.UnixNano()))params.Set(end,fmt.Sprintf(%d,end.UnixNano()))params.Set(limit,fmt.Sprintf(%d,limit))reqURL:fmt.Sprintf(%s/loki/api/v1/query_range?%s,c.baseURL,params.Encode())req,_:http.NewRequestWithContext(ctx,GET,reqURL,nil)req.Header.Set(Accept,application/json)resp,err:c.httpClient.Do(req)iferr!nil{returnnil,err}deferresp.Body.Close()// 解析Loki返回的流式数据varresultstruct{Datastruct{Result[]struct{Values[][]stringjson:values// [timestamp, log_line]}json:result}json:data}json.NewDecoder(resp.Body).Decode(result)// 提取日志条目varentries[]LogEntryfor_,stream:rangeresult.Data.Result{for_,v:rangestream.Values{iflen(v)2{entriesappend(entries,LogEntry{v[0],v[1]})}}}returnentries,nil}// SearchByTraceID 按trace_id搜索日志// 串联一次请求经过的所有服务日志func(c*LokiClient)SearchByTraceID(ctx context.Context,traceIDstring,lookback time.Duration)([]LogEntry,error){// LogQL查询: 所有服务中包含该trace_id的日志query:fmt.Sprintf({service~.} | trace_id%s,traceID)end:time.Now()start:end.Add(-lookback)returnc.QueryRange(ctx,query,start,end,5000)}// ESClient Elasticsearch客户端(简化版)typeESClientstruct{baseURLstringhttpClient*http.Client}// SearchLogs ES全文检索// 用query_string语法查询支持Kibana Query Languagefunc(c*ESClient)SearchLogs(ctx context.Context,index,querystring,from,sizeint)(map[string]interface{},error){// 构造ES查询体body:fmt.Sprintf({ query: {query_string: {query: %s}}, from: %d, size: %d, sort: [{timestamp: desc}] },query,from,size)reqURL:fmt.Sprintf(%s/%s/_search,c.baseURL,index)req,_:http.NewRequestWithContext(ctx,POST,reqURL,strings.NewReader(body))req.Header.Set(Content-Type,application/json)resp,err:c.httpClient.Do(req)iferr!nil{returnnil,err}deferresp.Body.Close()result:make(map[string]interface{})json.NewDecoder(resp.Body).Decode(result)returnresult,nil}// LogEntry 日志条目typeLogEntrystruct{Timestampstringjson:timestampLogstringjson:log}实际排查问题时先用trace_id搜索。Kibana里搜trace_id: abc123或Loki里搜{service~.} | trace_idabc123能看到一次请求从网关到业务服务到数据库的全部日志。然后按时间排序还原请求的执行过程。四、独家踩坑:日志量过大导致ES内存溢出这个坑很经典。30个服务每天50GB日志ES集群3个节点各配了16GB内存(ES推荐的JVM上限)。日志索引没有做生命周期管理历史日志一直堆积。第45天ES节点开始频繁Full GC查询超时最后OOM崩溃。根因有三个。第一所有日志存一个索引索引太大查询慢。第二没有设索引生命周期日志无限堆积。第三Debug日志没过滤大量无用日志挤占了存储。packagemainimport(fmtnet/httpstringstime)// CreateILMPolicy 创建ES索引生命周期策略// 通过ES API设置自动滚动和删除funcCreateILMPolicy(esURL,policyName,hotSizestring,deleteDaysint)error{// ILM策略: 热阶段达到hotSize滚动delete阶段按天数自动删除body:fmt.Sprintf({ policy: { phases: { hot: {actions: {rollover: {max_size: %s, max_age: 1d}}}, delete: {min_age: %dd, actions: {delete: {}}} } } },hotSize,deleteDays)req,_:http.NewRequest(PUT,fmt.Sprintf(%s/_ilm/policy/%s,esURL,policyName),strings.NewReader(body))req.Header.Set(Content-Type,application/json)client:http.Client{Timeout:10*time.Second}resp,err:client.Do(req)iferr!nil{returnerr}deferresp.Body.Close()returnnil}// ShouldLog 生产环境过滤Debug日志减少日志量funcShouldLog(levelstring)bool{switchlevel{caseerror,warn,info:returntruecasedebug,trace:returnfalse// 生产环境过滤default:returntrue}}修复后做了三件事。第一按服务名和日期建索引go-logs-order-service-2026.08.17。第二开ILM自动管理索引达到50GB或1天就滚动30天自动删除。第三生产环境过滤Debug日志日志量从50GB降到15GB。ES集群稳定运行不再OOM。如果日志量持续增长考虑引入Loki替代部分ES场景。Loki不索引正文只索引标签存储成本只有ES的十分之一。把全量日志存Loki只把Error日志存ES做全文检索兼顾成本和检索速度。五、对比分析维度ELKLoki索引方式全文倒排索引只索引标签存储成本高(10x)低(1x)全文检索快(毫秒级)慢(秒级顺序扫描)查询语言KQL/DSLLogQL部署复杂度高(ES集群LogstashKibana)低(LokiPromtailGrafana)适用场景需要全文检索按标签过滤日志查看ELK适合需要频繁全文检索的场景比如搜索某个错误关键字在所有服务中的出现。Loki适合按标签过滤查看日志的场景比如查看某个服务的Error日志。两者不冲突可以并存。全量日志存Loki关键日志存ES。实际项目里Loki覆盖80%的查询需求ES只处理20%的全文检索。总结日志采集的第一步是让zap输出结构化JSON字段名用Filebeat能识别的标准名(timestamp、message)。trace_id串联跨服务日志是排查问题的关键。ES必须配置ILM索引生命周期管理否则日志堆积会导致OOM。日志量大时引入Loki只索引标签不索引正文存储成本降一个数量级。上一篇讲了zap结构化日志这篇把日志采集和检索讲完日志体系从产出、采集到查询形成完整流程。