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会把请求体也复制:对大请求体(文件上传)要谨慎,可能占用大量内存。 - 金丝雀的回滚是"流量回滚"而不是"实例回滚":金丝雀实例依然在运行(用于排查)——需要确保它不再接收生产流量。