文档目录

8.6 配套代码:流量回放、副作用阻断与金丝雀

对应小节:8.6 流量回放与金丝雀发布 四部分:① 日志脱敏;② 副作用阻断;③ 影子流量配置;④ 金丝雀监控。

一、日志脱敏

# tools/sanitize_logs.py
"""
脱敏生产日志,准备用于流量回放。

两个原则:
  ① 保持格式(否则测不到校验逻辑)
  ② 保持关联(同一用户的请求要映射到同一新 ID,否则破坏缓存命中率)
"""
import hashlib
import json
import sys

# 需要脱敏的字段
ID_FIELDS = {"user_id", "order_id", "session_id", "account_id", "customer_id"}
PHONE_FIELDS = {"phone", "mobile", "phone_number"}
EMAIL_FIELDS = {"email", "user_email"}
TOKEN_FIELDS = {"token", "authorization", "api_key", "access_token", "refresh_token"}
NAME_FIELDS = {"name", "user_name", "real_name", "full_name"}


class Sanitizer:
    def __init__(self, salt="replay-salt-v1"):
        self.salt = salt
        self.id_map = {}
        self.stats = {"id": 0, "phone": 0, "email": 0, "token": 0, "name": 0}

    def stable_id(self, original: str, prefix: str = "replay") -> str:
        """保持关联的映射:同一输入 → 同一输出"""
        key = (prefix, original)
        if key not in self.id_map:
            h = hashlib.sha256((self.salt + prefix + original).encode()).hexdigest()[:12]
            self.id_map[key] = f"{prefix}-{h}"
        return self.id_map[key]

    def phone(self, value: str) -> str:
        """保持格式:11 位数字 → 11 位数字(避免影响校验逻辑)"""
        digits = "".join(c for c in str(value) if c.isdigit())
        if len(digits) == 11:
            return f"139{digits[-8:].rjust(8, '0')}"
        return "13900000000"

    def email(self, value: str) -> str:
        return "user@example.com"

    def name(self, value: str) -> str:
        return "测试用户"

    def token(self, value: str) -> str:
        return "replay-token-placeholder"

    def sanitize(self, obj):
        if isinstance(obj, dict):
            out = {}
            for k, v in obj.items():
                kl = k.lower()
                if kl in ID_FIELDS and v is not None:
                    out[k] = self.stable_id(str(v), prefix=kl.replace("_id", ""))
                    self.stats["id"] += 1
                elif kl in PHONE_FIELDS and v is not None:
                    out[k] = self.phone(str(v))
                    self.stats["phone"] += 1
                elif kl in EMAIL_FIELDS and v is not None:
                    out[k] = self.email(str(v))
                    self.stats["email"] += 1
                elif kl in TOKEN_FIELDS and v is not None:
                    out[k] = self.token(str(v))
                    self.stats["token"] += 1
                elif kl in NAME_FIELDS and v is not None:
                    out[k] = self.name(str(v))
                    self.stats["name"] += 1
                else:
                    out[k] = self.sanitize(v)
            return out
        elif isinstance(obj, list):
            return [self.sanitize(x) for x in obj]
        return obj


def main(in_path, out_path):
    s = Sanitizer()
    n_lines = 0

    with open(in_path, encoding="utf-8", errors="ignore") as fin, \
         open(out_path, "w", encoding="utf-8") as fout:
        for line in fin:
            if not line.strip():
                continue
            try:
                obj = json.loads(line)
            except json.JSONDecodeError:
                continue        # 跳过无法解析的行
            fout.write(json.dumps(s.sanitize(obj), ensure_ascii=False) + "\n")
            n_lines += 1

    print("═" * 70)
    print("脱敏完成")
    print("═" * 70)
    print(f"  输入: {in_path}")
    print(f"  输出: {out_path}")
    print(f"  处理行数: {n_lines}")
    print(f"  ID 映射数: {len(s.id_map)}(保持关联关系)")
    print()
    print("  字段统计:")
    for k, v in s.stats.items():
        print(f"    {k:<8}: {v}")
    print()
    print("⚠️  回放前必须人工检查:")
    print("  ① 随机抽样 20 行,确认没有遗漏的敏感字段")
    print("  ② 确认 ID 映射是稳定的(同输入 → 同输出)")
    print("  ③ 确认格式没变(手机号还是 11 位、邮箱还有 @)")
    print("═" * 70)


if __name__ == "__main__":
    if len(sys.argv) < 3:
        print("usage: sanitize_logs.py <input.jsonl> <output.jsonl>")
        sys.exit(1)
    main(sys.argv[1], sys.argv[2])

二、副作用阻断

// src/main/kotlin/safety/SideEffectGuard.kt
package safety

import io.micrometer.core.instrument.Counter
import io.micrometer.core.instrument.MeterRegistry
import org.slf4j.LoggerFactory

/**
 * 副作用阻断开关。
 *
 * 用途:在流量回放/影子环境中,阻止短信、邮件、支付、推送等真实副作用。
 *
 * ⚠️ 关键设计:
 *   ① 默认值必须是「启用副作用」——安全的默认是"不改变现有行为"
 *   ② 被阻断时必须【记录指标】——否则你不知道开关是否生效
 *   ③ 支持按类型精细控制(比如只允许邮件、阻断支付)
 */
object SideEffectGuard {
    private val log = LoggerFactory.getLogger(SideEffectGuard::class.java)
    private lateinit var registry: MeterRegistry

    /** 总开关:影子环境必须设为 false */
    @Volatile
    var enabled: Boolean = System.getenv("ENABLE_SIDE_EFFECTS")?.toBoolean() ?: true
        private set

    /** 精细控制:允许的副作用类型(空集合表示全部允许) */
    @Volatile
    private var allowedTypes: Set<String> = emptySet()

    private val blockedCounters = mutableMapOf<String, Counter>()

    fun init(registry: MeterRegistry) {
        this.registry = registry
    }

    fun configure(enabled: Boolean, allowed: Set<String> = emptySet()) {
        this.enabled = enabled
        this.allowedTypes = allowed
        log.warn("副作用开关已变更:enabled={}, allowed={}", enabled, allowed)
    }

    /**
     * 包裹一个有副作用的操作。
     *
     * @param type 副作用类型(sms / email / payment / push / external_api)
     * @return 被阻断时返回 null;否则返回操作结果
     */
    fun <T> guarded(type: String, description: String, block: () -> T): T? {
        val allowed = enabled && (allowedTypes.isEmpty() || type in allowedTypes)
        if (!allowed) {
            val counter = blockedCounters.computeIfAbsent(type) {
                Counter.builder("side_effect.blocked")
                    .tag("type", type)
                    .description("被阻断的副作用次数")
                    .register(registry)
            }
            counter.increment()
            log.warn("⛔ 副作用已阻断 [{}] {}(影子环境)", type, description)
            return null
        }
        return block()
    }

    /** 熔断式的保护:如果被阻断的比例过高,抛出异常(防止逻辑依赖副作用结果) */
    fun requireEnabled(type: String) {
        if (!enabled || (allowedTypes.isNotEmpty() && type !in allowedTypes)) {
            throw IllegalStateException("副作用 [$type] 在影子环境中被禁用,但代码依赖它的结果")
        }
    }
}

在真实业务代码里使用:

class OrderService(
    private val orderRepo: OrderRepository,
    private val smsClient: SmsClient,
    private val paymentClient: PaymentClient,
) {
    suspend fun createOrder(request: CreateOrderRequest): Order {
        val order = orderRepo.insert(request)

        // 短信:影子环境阻断
        SideEffectGuard.guarded("sms", "订单创建通知 ${order.id}") {
            smsClient.send(request.phone, "您的订单 ${order.id} 已创建")
        }

        // 支付:影子环境阻断(更需要谨慎)
        SideEffectGuard.guarded("payment", "扣款 ${order.id}") {
            paymentClient.charge(request.userId, order.amount)
        }

        return order
    }
}

启动配置(影子环境):

# 影子环境:阻断所有副作用
ENABLE_SIDE_EFFECTS=false java -jar shadow-app.jar

# 或者只允许邮件、阻断支付与短信
# 通过 SideEffectGuard.configure(false, setOf("email")) 在代码里配置

验证开关生效:

# 被阻断的副作用计数(应该 > 0)
sum by (type) (increase(side_effect.blocked_total[5m]))

# 如果影子环境跑了一段时间但这里全是 0 → 开关没生效!

三、影子流量配置

# nginx-shadow.conf —— 用 Nginx 做流量镜像
upstream production {
    server prod-app:8080;
}

upstream shadow {
    server shadow-app:8080;
}

server {
    listen 80;

    location /api/ {
        # ① 正常转发给生产
        proxy_pass http://production;

        # ② 同时镜像一份到影子环境
        mirror /shadow-mirror;
        mirror_request_body on;

        # ③ 限制镜像的超时(避免影子拖累生产)
        proxy_connect_timeout 1s;
        proxy_read_timeout 30s;
    }

    # 影子接收端(internal 表示外部无法直接访问)
    location = /shadow-mirror {
        internal;

        # ⚠️ 影子环境的响应会被丢弃,但不能让它无限等待
        proxy_pass http://shadow$request_uri;
        proxy_connect_timeout 1s;
        proxy_read_timeout 10s;

        # ④ 限制并发(避免影子环境被打爆而影响 Nginx)
        limit_conn shadow_limit 50;
    }
}

Istio 的等价配置:

apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
  name: app-mirror
spec:
  hosts: [app]
  http:
    - route:
        - destination:
            host: app-production
      mirror:
        host: app-shadow
      mirrorPercentage:
        value: 100.0        # 镜像 100% 流量(也可以只镜像 10%)

必须做的三件事:

事项 做法
影子环境阻断写操作 影子实例连影子数据库
阻断副作用 ENABLE_SIDE_EFFECTS=false
限制下游压力 影子环境的下游调用走 mock 或限流
# 影子环境的配置示例
env:
  - name: ENABLE_SIDE_EFFECTS
    value: "false"
  - name: DATABASE_URL
    value: "jdbc:postgresql://shadow-db:5432/app"     # 影子库
  - name: REDIS_URL
    value: "redis://shadow-redis:6379"                 # 影子缓存
  - name: DOWNSTREAM_USER_SERVICE
    value: "http://user-service-mock:8080"             # mock 下游
  - name: DOWNSTREAM_RATE_LIMIT_QPS
    value: "50"                                        # 限制对下游的压力

四、金丝雀监控与自动回滚

已在教学正文 8.6 节给出完整脚本。这里补充回滚脚本:

#!/usr/bin/env bash
# scripts/rollback-canary.sh
#
# 回滚金丝雀发布。
set -euo pipefail

echo "⛔ 开始回滚金丝雀"

# ① 记录回滚原因(供事后分析)
{
  echo "rollback_time=$(date -Iseconds)"
  echo "reason=${ROLLBACK_REASON:-未指定}"
  echo "canary_version=$(kubectl get deploy app-canary -o jsonpath='{.spec.template.spec.containers[0].image}' 2>/dev/null || echo n/a)"
} | tee "docs/experiments/rollback-$(date +%s).txt"

# ② 把流量全部切回稳定版本
kubectl patch service app -p '{"spec":{"selector":{"version":"stable"}}}'

# ③ 或按比例回退(如果用了 Istio)
# kubectl apply -f - <<EOF
# apiVersion: networking.istio.io/v1beta1
# kind: VirtualService
# metadata:
#   name: app
# spec:
#   hosts: [app]
#   http:
#     - route:
#         - destination: { host: app-stable }
#           weight: 100
# EOF

# ④ 确认回滚生效
sleep 5
echo "当前流量分配:"
kubectl get virtualservice app -o jsonpath='{.spec.http[0].route}' | python3 -m json.tool

# ⑤ 保留金丝雀实例(不要立即删除,以便排查)
echo
echo "⚠️  金丝雀实例保留(未删除),用于排查:"
echo "   kubectl logs deploy/app-canary --tail=1000 > canary-logs.txt"
echo "   kubectl describe pod -l version=canary > canary-describe.txt"

echo
echo "✅ 回滚完成"
echo
echo "后续:"
echo "  ① 收集金丝雀的证据(日志、指标、火焰图)"
echo "  ② 按第 6 章的流程定位根因"
echo "  ③ 修复后重新走金丝雀流程"

五、动手改造

改动 观察什么
用 sanitize_logs.py 处理一份真实日志 抽样检查有没有遗漏的敏感字段
在影子环境不设 ENABLE_SIDE_EFFECTS=false 会真的发短信/扣款——这就是为什么必须显式验证开关
检查 side_effect.blocked_total 指标 确认开关真的生效(指标为 0 说明没生效)
给影子环境的 Nginx 去掉 limit_conn 影子可能打垮 Nginx——理解"必须限流"
用金丝雀监控跑一次真实发布 观察它能否检出问题并触发回滚

六、这段代码的局限

  • sanitize_logs.py 只能处理 JSON 格式的日志:如果是文本格式,需要改成正则替换。
  • 脱敏不保证 100% 完整:可能有嵌套在字符串里的敏感信息(比如 URL 参数)——必须人工抽样检查。
  • SideEffectGuard 需要改业务代码:如果代码里有很多散落的副作用调用,接入成本高(但这是必须付的成本)。
  • Nginx 的 mirror 会把请求体也复制:对大请求体(文件上传)要谨慎,可能占用大量内存。
  • 金丝雀的回滚是"流量回滚"而不是"实例回滚":金丝雀实例依然在运行(用于排查)——需要确保它不再接收生产流量。