ELK 日志收集 + 分布式追踪高级配置 + 全链路安全加固

2025-01-04T14:21:02+08:00 | 10分钟阅读 | 更新于 2025-01-04T14:21:02+08:00

@

学习目标

类比(先建立直觉):把生产环境想象成一座大型工厂。ELK 是工厂的"监控录像 + 报警大屏"(日志收集与检索),OpenTelemetry 是"全流程工单追踪"(一次请求在哪些车间流转),全链路安全加固则是"门禁 + 车间互锁 + 保险柜"(HTTPS / mTLS / Vault)。本章把这三套系统整合起来,让一个微服务既能"被看见"、又能"被追踪"、还能"防住坏人"。

学完本章你应该能够:

  1. 在 K8s 上部署一套 Filebeat + Elasticsearch + Kibana 日志链路,并让应用日志自动带上 trace_id,实现"按 TraceID 查全链路日志"。
  2. 在 OpenTelemetry 中配置 自定义属性、事件、智能采样(错误 100% + 正常按比例)与 Baggage 跨服务透传
  3. Istio mTLS、HTTPS 强制跳转、CORS、速率限制 给 API 网关与服务间通信做安全加固。
  4. Vault 管理数据库密码等敏感信息,并用 非 root 容器 + SecurityContext + NetworkPolicy 收紧 K8s 运行时安全。
  5. 把上面所有能力拼成一张"生产级微服务全景图",并在面试中讲清每一层的取舍。

前置知识

  • K8s 基础:Namespace / Deployment / StatefulSet / DaemonSet / ConfigMap / Service
  • Go + Kitex/Hertz 基础,了解 context.Context 与中间件写法
  • OpenTelemetry 基础概念(见上一篇《Golang 接入 OpenTelemetry》)
  • 基本的 HTTP 安全常识(HTTPS、JWT、CORS)

本章你会动手做的事

  1. 按文中 YAML 在本地 K8s(或 Kind)起一套 ELK,把示例服务的日志灌进去并在 Kibana 用 trace_id 检索。
  2. 给订单创建接口加自定义属性与事件,并把采样策略改成"错误全采、正常 10%"。
  3. kubectl 给某服务加上 SecurityContext + NetworkPolicy,验证非 root 与 Pod 间网络隔离生效。

一、ELK 全链路日志收集(Filebeat+Elasticsearch+Kibana)

字节跳动内部日志系统核心架构:应用输出结构化 JSON 日志 → Filebeat 采集 → Kafka 缓冲 → Logstash 过滤 → Elasticsearch 存储 → Kibana 可视化,支持 PB 级日志秒级检索

这张图在讲什么:一条应用日志从产生到可被检索,依次经过采集、缓冲、过滤、存储、可视化五个环节。Kafka 在这里承担"削峰填谷"——日志洪峰时先堆在 Kafka,Logstash 按自己节奏消费,避免 ES 被冲垮。

flowchart LR
    App[应用 Pod
结构化 JSON 日志] --> FB[Filebeat
采集] FB --> KF[Kafka
缓冲削峰] KF --> LS[Logstash
过滤/富化] LS --> ES[Elasticsearch
存储/检索] ES --> Kib[Kibana
可视化]

1. 完整 ELK Stack 部署

# k8s/elk/namespace.yaml
apiVersion: v1
kind: Namespace
metadata:
  name: elk
---
# k8s/elk/elasticsearch.yaml
apiVersion: apps/v1
kind: StatefulSet
metadata:
  name: elasticsearch
  namespace: elk
spec:
  replicas: 3
  selector:
    matchLabels:
      app: elasticsearch
  template:
    metadata:
      labels:
        app: elasticsearch
    spec:
      containers:
      - name: elasticsearch
        image: docker.elastic.co/elasticsearch/elasticsearch:8.13.0
        ports:
        - containerPort: 9200
        - containerPort: 9300
        env:
        - name: discovery.type
          value: zen
        - name: ES_JAVA_OPTS
          value: "-Xms2g -Xmx2g"
        - name: ELASTIC_PASSWORD
          value: "elastic123"
        volumeMounts:
        - name: es-data
          mountPath: /usr/share/elasticsearch/data
  volumeClaimTemplates:
  - metadata:
      name: es-data
    spec:
      accessModes: ["ReadWriteOnce"]
      resources:
        requests:
          storage: 100Gi
---
apiVersion: v1
kind: Service
metadata:
  name: elasticsearch
  namespace: elk
spec:
  selector:
    app: elasticsearch
  ports:
  - port: 9200
    targetPort: 9200
---
# k8s/elk/kibana.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: kibana
  namespace: elk
spec:
  replicas: 1
  selector:
    matchLabels:
      app: kibana
  template:
    metadata:
      labels:
        app: kibana
    spec:
      containers:
      - name: kibana
        image: docker.elastic.co/kibana/kibana:8.13.0
        ports:
        - containerPort: 5601
        env:
        - name: ELASTICSEARCH_HOSTS
          value: "http://elasticsearch:9200"
        - name: ELASTICSEARCH_USERNAME
          value: "elastic"
        - name: ELASTICSEARCH_PASSWORD
          value: "elastic123"
---
apiVersion: v1
kind: Service
metadata:
  name: kibana
  namespace: elk
spec:
  type: LoadBalancer
  selector:
    app: kibana
  ports:
  - port: 5601
    targetPort: 5601
---
# k8s/elk/filebeat.yaml
apiVersion: apps/v1
kind: DaemonSet
metadata:
  name: filebeat
  namespace: elk
spec:
  selector:
    matchLabels:
      app: filebeat
  template:
    metadata:
      labels:
        app: filebeat
    spec:
      containers:
      - name: filebeat
        image: docker.elastic.co/beats/filebeat:8.13.0
        volumeMounts:
        - name: filebeat-config
          mountPath: /usr/share/filebeat/filebeat.yml
          subPath: filebeat.yml
        - name: varlog
          mountPath: /var/log
        - name: varlibdockercontainers
          mountPath: /var/lib/docker/containers
          readOnly: true
      volumes:
      - name: filebeat-config
        configMap:
          name: filebeat-config
      - name: varlog
        hostPath:
          path: /var/log
      - name: varlibdockercontainers
        hostPath:
          path: /var/lib/docker/containers
---
apiVersion: v1
kind: ConfigMap
metadata:
  name: filebeat-config
  namespace: elk
data:
  filebeat.yml: |
    filebeat.inputs:
    - type: container
      paths:
        - /var/log/containers/*.log
      processors:
        - add_kubernetes_metadata:
            host: ${NODE_NAME}
            matchers:
            - logs_path:
                logs_path: "/var/log/containers/"

    output.elasticsearch:
      hosts: ["elasticsearch:9200"]
      username: "elastic"
      password: "elastic123"
      index: "mall-logs-%{+yyyy.MM.dd}"

    setup.ilm.enabled: false
    setup.template.name: "mall-logs"
    setup.template.pattern: "mall-logs-*"    

2. 应用日志增强(集成 TraceID)

修改pkg/logger/logger.go,自动在日志中注入链路追踪 ID:

package logger

import (
    "context"
    "os"
    "go.opentelemetry.io/otel/trace"
    "go.uber.org/zap"
    "go.uber.org/zap/zapcore"
)

var Logger *zap.Logger

// Ctx 从上下文获取trace_id和span_id,返回带上下文的logger
func Ctx(ctx context.Context) *zap.Logger {
    span := trace.SpanFromContext(ctx)
    if span.SpanContext().IsValid() {
        return Logger.With(
            zap.String("trace_id", span.SpanContext().TraceID().String()),
            zap.String("span_id", span.SpanContext().SpanID().String()),
        )
    }
    return Logger
}

// 所有日志方法都支持上下文参数
func InfoCtx(ctx context.Context, msg string, fields ...zap.Field) {
    Ctx(ctx).Info(msg, fields...)
}

func ErrorCtx(ctx context.Context, msg string, fields ...zap.Field) {
    Ctx(ctx).Error(msg, fields...)
}

func FatalCtx(ctx context.Context, msg string, fields ...zap.Field) {
    Ctx(ctx).Fatal(msg, fields...)
}

3. 业务代码中使用增强日志

// 示例:internal/userservice/service/user_service.go
func (s *UserServiceImpl) GetUserInfo(ctx context.Context, req *user.GetUserInfoRequest) (*user.GetUserInfoResponse, error) {
    // 自动包含trace_id和span_id
    logger.InfoCtx(ctx, "GetUserInfo request", zap.Uint64("user_id", req.UserId))

    userModel, err := s.userDAO.GetByID(ctx, req.UserId)
    if err != nil {
        logger.ErrorCtx(ctx, "GetUserInfo failed", zap.Error(err))
        return nil, err
    }

    return &user.GetUserInfoResponse{
        UserId:   userModel.ID,
        Username: userModel.Username,
        Email:    userModel.Email,
        Phone:    userModel.Phone,
    }, nil
}

4. Kibana 日志查询最佳实践

  • 按 TraceID 查询全链路日志trace_id:"xxxxxx"
  • 按服务查询kubernetes.labels.app:"userservice"
  • 按错误级别查询level:"ERROR"
  • 按时间范围查询@timestamp:[now-1h TO now]

二、OpenTelemetry 分布式追踪高级配置

1. 自定义追踪属性与事件

// 示例:在订单创建接口添加自定义追踪信息
func (s *OrderServiceImpl) CreateOrder(ctx context.Context, req *order.CreateOrderRequest) (*order.CreateOrderResponse, error) {
    userID := ctx.Value("user_id").(uint64)
    
    // 获取当前span
    span := trace.SpanFromContext(ctx)
    
    // 添加自定义属性
    span.SetAttributes(
        attribute.Int64("user_id", int64(userID)),
        attribute.Int("item_count", len(req.Items)),
        attribute.Float64("total_amount", totalAmount),
    )
    
    // 添加事件
    span.AddEvent("开始创建订单")
    
    // ... 业务逻辑
    
    span.AddEvent("订单创建成功", trace.WithAttributes(
        attribute.Int64("order_id", int64(orderModel.ID)),
    ))
    
    return &order.CreateOrderResponse{OrderId: int64(orderModel.ID)}, nil
}

2. 智能采样策略(高并发必备)

// pkg/otel/otel.go
func InitProvider(serviceName, endpoint string) provider.OtelProvider {
    // 混合采样器:
    // 1. 所有错误请求100%采样
    // 2. 正常请求按10%概率采样
    sampler := sdktrace.NewParentBasedSampler(
        sdktrace.NewTraceIDRatioBased(0.1),
        sdktrace.WithRemoteParentSampled(sdktrace.AlwaysSample()),
        sdktrace.WithRemoteParentNotSampled(sdktrace.NeverSample()),
    )

    return provider.NewOpenTelemetryProvider(
        provider.WithServiceName(serviceName),
        provider.WithExportEndpoint(endpoint),
        provider.WithInsecure(),
        provider.WithSampler(sampler), // 应用采样策略
    )
}

3. Baggage 跨服务上下文透传

用于在整个调用链中传递用户 ID、请求 ID 等全局信息:

这张图在讲什么:网关在入口把 user_id 放进 Baggage,OTel 的 propagator 会自动把它透传到下游每个服务,无需在业务参数里层层手动传递。

flowchart LR
    G[网关 JWTAuth
设置 Baggage user_id] --> A[订单服务] A -->|自动透传| B[用户服务] A -->|自动透传| C[支付服务]
// 网关层设置Baggage
func JWTAuth() app.HandlerFunc {
    return func(ctx context.Context, c *app.RequestContext) {
        // ... 验证token
        
        // 将用户ID添加到Baggage,自动透传到所有下游服务
        baggageMember, _ := baggage.NewMember("user_id", strconv.FormatUint(claims.UserID, 10))
        baggage, _ := baggage.New(baggageMember)
        ctx = baggage.ContextWithBaggage(ctx, baggage)
        
        c.Next(ctx)
    }
}

// 下游服务获取Baggage
func (s *OrderServiceImpl) CreateOrder(ctx context.Context, req *order.CreateOrderRequest) (*order.CreateOrderResponse, error) {
    // 从Baggage获取用户ID
    b := baggage.FromContext(ctx)
    userIDStr := b.Member("user_id").Value()
    userID, _ := strconv.ParseUint(userIDStr, 10, 64)
    
    // ... 业务逻辑
}

4. 日志与链路追踪关联

在 Kibana 中点击日志中的trace_id,可以直接跳转到 Jaeger 查看完整链路,实现日志 - 链路一键跳转。

这张图在讲什么:一条带 trace_id 的日志,从 Elasticsearch 被 Kibana 展示,点击 trace_id 即可跳转到 Jaeger 看到整条调用链,并定位到具体失败的 Span。

flowchart LR
    L[ES 中的一条日志
带 trace_id] -->|点击 trace_id| K[Kibana] K -->|跳转| J[Jaeger
展示完整 Trace] J -->|定位失败 Span| S[具体服务/方法]

三、全链路安全加固

1. API 网关安全防护

1.1 HTTPS 强制加密

# k8s/istio/gateway.yaml
apiVersion: networking.istio.io/v1alpha3
kind: Gateway
metadata:
  name: mall-gateway
  namespace: cloudwego-mall
spec:
  selector:
    istio: ingressgateway
  servers:
  - port:
      number: 80
      name: http
      protocol: HTTP
    hosts:
    - "mall.example.com"
    tls:
      httpsRedirect: true # 强制HTTP跳转到HTTPS
  - port:
      number: 443
      name: https
      protocol: HTTPS
    hosts:
    - "mall.example.com"
    tls:
      mode: SIMPLE
      credentialName: mall-tls # TLS证书

1.2 CORS 跨域安全配置

// internal/api-gateway/middleware/cors.go
package middleware

import (
    "context"
    "github.com/cloudwego/hertz/pkg/app"
)

func CORS() app.HandlerFunc {
    return func(ctx context.Context, c *app.RequestContext) {
        c.Header("Access-Control-Allow-Origin", "https://mall.example.com")
        c.Header("Access-Control-Allow-Methods", "GET, POST, PUT, DELETE, OPTIONS")
        c.Header("Access-Control-Allow-Headers", "Content-Type, Authorization")
        c.Header("Access-Control-Max-Age", "86400")
        c.Header("X-Content-Type-Options", "nosniff")
        c.Header("X-Frame-Options", "DENY")
        c.Header("X-XSS-Protection", "1; mode=block")

        if c.Request.Method() == "OPTIONS" {
            c.AbortWithStatus(204)
            return
        }

        c.Next(ctx)
    }
}

1.3 接口速率限制

// internal/api-gateway/middleware/ratelimit.go
package middleware

import (
    "context"
    "time"
    "github.com/cloudwego/hertz/pkg/app"
    "github.com/redis/go-redis/v9"
    "github.com/yourname/cloudwego-mall/pkg/redis"
    "github.com/yourname/cloudwego-mall/pkg/errors"
    "github.com/yourname/cloudwego-mall/internal/api-gateway/model"
)

func RateLimit() app.HandlerFunc {
    return func(ctx context.Context, c *app.RequestContext) {
        // 按IP限流,每分钟最多100次请求
        ip := c.ClientIP()
        key := "ratelimit:" + ip

        count, err := redis.Client.Incr(ctx, key).Result()
        if err != nil {
            model.Error(c, errors.New(errors.CodeInternalError, "限流服务异常"))
            c.Abort()
            return
        }

        if count == 1 {
            redis.Client.Expire(ctx, key, 1*time.Minute)
        }

        if count > 100 {
            model.Error(c, errors.New(errors.CodeRateLimitExceeded, "请求过于频繁,请稍后再试"))
            c.Abort()
            return
        }

        c.Next(ctx)
    }
}

2. 微服务间 mTLS 双向认证

使用 Istio 实现服务间自动双向 TLS 加密,防止中间人攻击

# k8s/istio/peer-authentication.yaml
apiVersion: security.istio.io/v1beta1
kind: PeerAuthentication
metadata:
  name: default
  namespace: cloudwego-mall
spec:
  mtls:
    mode: STRICT # 强制所有服务间通信使用mTLS

3. 敏感信息加密与配置安全

3.1 使用 Vault 管理敏感信息

不要在配置文件中硬编码数据库密码、API 密钥等敏感信息,使用 HashiCorp Vault 统一管理:

// pkg/vault/vault.go
package vault

import (
    "github.com/hashicorp/vault/api"
    "github.com/yourname/cloudwego-mall/pkg/config"
)

var Client *api.Client

func Init() error {
    cfg := api.DefaultConfig()
    cfg.Address = config.GlobalConfig.Vault.Address

    var err error
    Client, err = api.NewClient(cfg)
    if err != nil {
        return err
    }

    Client.SetToken(config.GlobalConfig.Vault.Token)
    return nil
}

// GetSecret 获取敏感信息
func GetSecret(path string) (map[string]interface{}, error) {
    secret, err := Client.Logical().Read(path)
    if err != nil {
        return nil, err
    }
    return secret.Data, nil
}

3.2 数据库密码加密示例

// pkg/mysql/mysql.go
func Init() error {
    // 从Vault获取数据库密码
    secret, err := vault.GetSecret("database/mysql")
    if err != nil {
        return err
    }
    password := secret["password"].(string)

    dsn := fmt.Sprintf("%s:%s@tcp(%s)/%s?charset=utf8mb4&parseTime=True&loc=Local",
        config.GlobalConfig.Mysql.Username,
        password,
        config.GlobalConfig.Mysql.Address,
        config.GlobalConfig.Mysql.Database,
    )

    // ... 剩余代码
}

4. 容器与 Kubernetes 安全加固

4.1 非 root 用户运行容器

# Dockerfile 示例
FROM golang:1.22-alpine AS builder
WORKDIR /app
COPY . .
RUN go build -o userservice cmd/userservice/main.go

FROM alpine:3.19
RUN addgroup -S appgroup && adduser -S appuser -G appgroup
USER appuser # 使用非root用户
WORKDIR /app
COPY --from=builder /app/userservice .
EXPOSE 8081
CMD ["./userservice"]

4.2 Kubernetes SecurityContext

# k8s/userservice.yaml 新增
spec:
  securityContext:
    runAsNonRoot: true
    runAsUser: 1000
    runAsGroup: 1000
    fsGroup: 1000
  containers:
  - name: userservice
    image: your-docker-registry/userservice:latest
    securityContext:
      allowPrivilegeEscalation: false
      readOnlyRootFilesystem: true # 只读根文件系统
      capabilities:
        drop:
        - ALL # 删除所有Linux能力
    volumeMounts:
    - name: tmp
      mountPath: /tmp
  volumes:
  - name: tmp
    emptyDir: {}

4.3 网络策略限制 Pod 间通信

# k8s/security/network-policy.yaml
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
  name: default-deny-all
  namespace: cloudwego-mall
spec:
  podSelector: {}
  policyTypes:
  - Ingress
  - Egress
---
# 允许API网关访问所有服务
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
  name: allow-api-gateway
  namespace: cloudwego-mall
spec:
  podSelector: {}
  ingress:
  - from:
    - podSelector:
        matchLabels:
          app: api-gateway
  policyTypes:
  - Ingress
---
# 允许订单服务访问商品和支付服务
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
  name: allow-order-service
  namespace: cloudwego-mall
spec:
  podSelector:
    matchLabels:
      app: orderservice
  egress:
  - to:
    - podSelector:
        matchLabels:
          app: productservice
    - podSelector:
        matchLabels:
          app: paymentservice
  policyTypes:
  - Egress

这张图在讲什么:一次外网请求要连过三道门——边缘 HTTPS 网关、API 网关的 CORS / 限流、以及服务间 mTLS,敏感配置再交给 Vault 统一管理。整体构成一个纵深防御(defense in depth)体系。

flowchart LR
    U[用户] -->|HTTPS 强制跳转| GW[API 网关
CORS/限流] GW -->|mTLS 双向认证| S1[服务A] S1 -->|mTLS| S2[服务B] S1 -.->|读密钥| V[Vault
敏感信息管理]

四、项目最终完整生产级能力全景

能力维度已实现特性
业务能力用户认证、购物车、商品管理、订单管理、支付集成
架构能力微服务架构、API 网关、服务注册发现、配置中心
高并发能力Redis 多级缓存、库存预扣减、异步消息队列、连接池优化
一致性能力Seata 分布式事务、乐观锁防超卖、最终一致性补偿
可观测性Zap 结构化日志、ELK 日志收集、OpenTelemetry 全链路追踪、Prometheus 监控、Grafana 可视化
高可用能力Sentinel 熔断限流、Istio 流量治理、金丝雀灰度发布、K8s 多副本部署、健康检查
安全能力HTTPS 加密、JWT 认证、mTLS 双向认证、CORS 防护、XSS/CSRF 防护、速率限制、敏感信息加密、容器安全、网络隔离
部署能力Docker 容器化、K8s 编排、完整部署脚本、CI/CD 集成

五、面试加分亮点总结

  1. 微服务架构设计:服务拆分原则、API 网关设计、服务间通信方式
  2. 高并发处理:缓存设计、库存防超卖、异步解耦、性能优化
  3. 分布式系统:分布式事务、分布式锁、服务注册发现、链路追踪
  4. 可观测性:日志、监控、告警、追踪的最佳实践
  5. 生产级安全:HTTPS、mTLS、敏感信息保护、容器安全
  6. 工程化能力:项目结构、代码规范、CI/CD、K8s 部署

自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. ELK 链路里 Kafka 的作用是什么?没有 Kafka(Filebeat 直连 Logstash)会有什么风险?
  2. 为什么应用日志里要注入 trace_id?在 Kibana 里按 trace_id 查询能解决什么排错痛点?
  3. OpenTelemetry 智能采样里"错误 100% + 正常 10%“的目的是什么?全采和全不采各有什么代价?
  4. Istio PeerAuthentication 设成 STRICT 意味着什么?它防的是哪一类攻击?
  5. SecurityContextrunAsNonRootreadOnlyRootFilesystem 分别防的是什么?NetworkPolicy 默认拒绝一切入站 / 出站后又单独放行网关,是为了解决什么?

动手练习(建议真做一遍)

  1. 在本地 Kind 集群按文中 YAML 起一套 ELK,造一条带 trace_id 的日志,验证 Kibana 能检索并在点击后跳到 Jaeger。
  2. 把订单服务的采样策略改成"错误全采、正常 10%",用压测制造一批正常 + 一个错误,观察后端 Trace 数量比例。
  3. 给某个服务加 NetworkPolicy 只允许来自 API 网关的入站,再用另一个 Pod curl 它,验证被拒绝,理解零信任网络。

本章小结

  • 日志可检索:Filebeat 采集 →(Kafka 缓冲)→ Logstash 过滤 → ES 存储 → Kibana 检索,是 PB 级日志链路的经典骨架。
  • 日志即链路:在日志里注入 trace_id、在 Kibana 点击跳转 Jaeger,实现"日志 — 链路一键定位”。
  • 追踪可增强:自定义 Attribute / Event、智能采样、Baggage 跨服务透传,是生产级 OTel 的三件套。
  • 边缘要收口:HTTPS 强制跳转、CORS、速率限制放在 API 网关,统一挡住外网常见风险。
  • 内部要互信:Istio mTLS STRICT 让服务间自动双向认证,防中间人;Vault 统一管理密钥,杜绝硬编码。
  • 运行时要最小权限:非 root 容器 + SecurityContext + NetworkPolicy 默认拒绝,构成 K8s 零信任底座。

至此,微服务从"能跑"走到了"可观测、可追踪、可防御"的生产级水准;后续可继续深入每一层的性能调优与多集群容灾。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!