学习目标
类比(先建立直觉):把生产环境想象成一座大型工厂。ELK 是工厂的"监控录像 + 报警大屏"(日志收集与检索),OpenTelemetry 是"全流程工单追踪"(一次请求在哪些车间流转),全链路安全加固则是"门禁 + 车间互锁 + 保险柜"(HTTPS / mTLS / Vault)。本章把这三套系统整合起来,让一个微服务既能"被看见"、又能"被追踪"、还能"防住坏人"。
学完本章你应该能够:
- 在 K8s 上部署一套 Filebeat + Elasticsearch + Kibana 日志链路,并让应用日志自动带上
trace_id,实现"按 TraceID 查全链路日志"。 - 在 OpenTelemetry 中配置 自定义属性、事件、智能采样(错误 100% + 正常按比例)与 Baggage 跨服务透传。
- 用 Istio mTLS、HTTPS 强制跳转、CORS、速率限制 给 API 网关与服务间通信做安全加固。
- 用 Vault 管理数据库密码等敏感信息,并用 非 root 容器 + SecurityContext + NetworkPolicy 收紧 K8s 运行时安全。
- 把上面所有能力拼成一张"生产级微服务全景图",并在面试中讲清每一层的取舍。
前置知识:
- K8s 基础:
Namespace/Deployment/StatefulSet/DaemonSet/ConfigMap/Service - Go + Kitex/Hertz 基础,了解
context.Context与中间件写法 - OpenTelemetry 基础概念(见上一篇《Golang 接入 OpenTelemetry》)
- 基本的 HTTP 安全常识(HTTPS、JWT、CORS)
本章你会动手做的事:
- 按文中 YAML 在本地 K8s(或 Kind)起一套 ELK,把示例服务的日志灌进去并在 Kibana 用
trace_id检索。 - 给订单创建接口加自定义属性与事件,并把采样策略改成"错误全采、正常 10%"。
- 用
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 集成 |
五、面试加分亮点总结
- 微服务架构设计:服务拆分原则、API 网关设计、服务间通信方式
- 高并发处理:缓存设计、库存防超卖、异步解耦、性能优化
- 分布式系统:分布式事务、分布式锁、服务注册发现、链路追踪
- 可观测性:日志、监控、告警、追踪的最佳实践
- 生产级安全:HTTPS、mTLS、敏感信息保护、容器安全
- 工程化能力:项目结构、代码规范、CI/CD、K8s 部署
自测题与动手练习
自测题(合上书能答出来,才算懂):
- ELK 链路里 Kafka 的作用是什么?没有 Kafka(Filebeat 直连 Logstash)会有什么风险?
- 为什么应用日志里要注入
trace_id?在 Kibana 里按trace_id查询能解决什么排错痛点? - OpenTelemetry 智能采样里"错误 100% + 正常 10%“的目的是什么?全采和全不采各有什么代价?
- Istio
PeerAuthentication设成STRICT意味着什么?它防的是哪一类攻击? SecurityContext里runAsNonRoot与readOnlyRootFilesystem分别防的是什么?NetworkPolicy默认拒绝一切入站 / 出站后又单独放行网关,是为了解决什么?
动手练习(建议真做一遍):
- 在本地 Kind 集群按文中 YAML 起一套 ELK,造一条带
trace_id的日志,验证 Kibana 能检索并在点击后跳到 Jaeger。 - 把订单服务的采样策略改成"错误全采、正常 10%",用压测制造一批正常 + 一个错误,观察后端 Trace 数量比例。
- 给某个服务加
NetworkPolicy只允许来自 API 网关的入站,再用另一个 Podcurl它,验证被拒绝,理解零信任网络。
本章小结
- 日志可检索: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 零信任底座。
至此,微服务从"能跑"走到了"可观测、可追踪、可防御"的生产级水准;后续可继续深入每一层的性能调优与多集群容灾。