模式与反模式
模式与反模式
钩子是一条小小的扩展缝,用好用坏都很容易。这一页收集在生产里站得住的形状,以及会咬人的那些。
模式
审计日志
记录每一个决策,并记下足以重建它的东西:状态、问题、答案、模型、路由决策、usage 和延迟。
import json
def audit(ctx):
for state, result in zip(ctx.states, ctx.results or []):
json.dump({
"run_id": ctx.run_id,
"model": ctx.model,
"state": state,
"routing": result.get("routing"),
"answers": result["answers"],
"usage": result.get("usage"),
"call_usage": ctx.usage,
"call_elapsed_ms": round(ctx.elapsed_ms or 0.0, 3),
}, sys.stdout)
sys.stdout.write("\n")
laya.load("convaiinnovations/laya", on_predict_end=audit)
一个钩子每次调用触发一次,而一次 predict_batch 调用把每个状态都装在里面,所以记录是按决策写
的:ctx.states 和 ctx.results 按下标对齐。ctx.usage 和 ctx.elapsed_ms 是整次调用的合计;
每个结果带它自己的 usage。
如果丢一行日志绝不能导致请求失败,就让它宽松:hooks_raise=False。如果审计轨迹是一项合规要求,
就让它严格。
PII 脱敏
脱敏必须发生在 on_predict_start 里、tokenization 之前,否则模型已经见过那些数据了。
import re
EMAIL = re.compile(r"\b[\w.+-]+@[\w-]+\.[\w.-]+\b")
def redact(ctx):
ctx.states = [
EMAIL.sub("[email]", s) if isinstance(s, str) else s
for s in ctx.states
]
laya.load("convaiinnovations/laya", on_predict_start=redact)
脱敏钩子是一条策略钩子:保持 hooks_raise=True,因为一个静默坏掉的脱敏器就是一次数据泄漏。
缓存
一个 start 钩子查缓存并调用 ctx.skip(...);一个 end 钩子填充它。命中时前向传播被跳过。
import hashlib, json
CACHE = {}
def key(ctx, index):
# Not sort_keys=True: criteria order is positional, so two orders are two questions,
# and the checkpoint and token budget change the answer too.
payload = json.dumps([ctx.states[index], ctx.questions, ctx.model,
ctx.max_len, ctx.head_max_len], default=str)
return hashlib.sha256(payload.encode()).hexdigest()
def read(ctx):
hits = [CACHE.get(key(ctx, i)) for i in range(len(ctx.states))]
if all(hit is not None for hit in hits):
ctx.skip(hits) # one per state: skip replaces the whole call
def write(ctx):
for i, result in enumerate(ctx.results or []):
CACHE[key(ctx, i)] = result
laya.load("convaiinnovations/laya", on_predict_start=read, on_predict_end=write)
钩子每次调用触发一次,所以在 predict_batch 上只按一个状态做 key 是不够的:ctx.skip() 会替换
这次调用本会返回的每一个结果。并发服务时用一把锁守住缓存。在 Router 上,被缓存的载荷仍然会拿到
一个 routing 键,所以返回形状不变。
同一对钩子也能通过单个 LangChain 节点的 hooks= 参数用在它上面,这是在不改变
该 agent 其他每个调用方所见内容的前提下,缓存图里某一步热点的办法。
指标
从 ctx.model、ctx.usage 和 ctx.elapsed_ms 得到计数器和直方图。让它宽松。
COUNTS, LATENCIES = {}, []
def metrics(ctx):
COUNTS[ctx.model] = COUNTS.get(ctx.model, 0) + 1
if ctx.elapsed_ms is not None:
LATENCIES.append(ctx.elapsed_ms)
laya.load("convaiinnovations/laya", on_predict_end=metrics, hooks_raise=False)
防护栏
一条策略钩子抛出以拦下一个请求。hooks_raise=True(默认)让拦截抵达调用方;on_error 和
on_predict_end 仍然运行,于是审计轨迹记下它。
class Blocked(Exception):
pass
def guard(ctx):
if any("ssn" in str(state).lower() for state in ctx.states):
raise Blocked("possible PII in state")
laya.load("convaiinnovations/laya", on_predict_start=guard)
拿状态形状去测它。一个只读 ctx.states[0] 的防护栏会拦下一次单调用,却让一次 predict_batch
调用把剩下的每一个状态都送进前向传播。
置信度门控
一个 end 钩子把低置信度的答案改写成安全的回退值,或者为下游逻辑给它加标注。这是一次结果改动,不是 一次拒绝。
def gate(ctx):
for result in ctx.results or []:
answer = result["answers"].get("dept")
if answer and answer["confidence"] < 0.6:
answer["choice"] = "human-review"
answer["gated"] = True
laya.load("convaiinnovations/laya", on_predict_end=gate)
通过 ctx.results 来改动,它为这次调用的每个状态保存一个 dict:只门控第一个,就会把其他每一个
低置信度的答案不加标注地发出去。
路由覆盖
on_route 可以替换 ctx.decision,为某一类流量钉住一个 checkpoint。
from laya.router import RouteDecision
def pin(ctx):
if "refund" in str(ctx.states[0]).lower():
ctx.decision = RouteDecision(
model="typed-decisions",
repo="convaiinnovations/laya/typed-decisions",
reason="refund workflow",
detection=None,
workflow=None,
)
Router(hooks=[pin])
模型生命周期
on_load 和 on_evict 观察 checkpoint。用它们做预热日志、内存核算或驱逐告警。它们运行在 Router
锁之外,所以钩子可以回调进 Router。
class Lifecycle:
def on_load(self, ctx):
print("loaded", ctx.model)
def on_evict(self, ctx):
print("evicted", ctx.model)
Router(hooks=[Lifecycle()])
多租户上下文
通过在钩子闭包里捕获它,或者从一个 context-local 里读它,来把租户 id 串起来。不要在钩子对象上存逐 请求的状态,除非上锁。
def make_audit(tenant):
def audit(ctx):
ship(tenant, ctx.run_id, ctx.results)
return audit
agent = laya.load("convaiinnovations/laya", on_predict_end=make_audit("acme"))
组合
几种不同种类的钩子自然地组合;已安装的钩子按顺序先运行。
agent = laya.load(
"convaiinnovations/laya",
hooks=[Metrics(), Guardrail()], # metrics first, then policy
on_predict_start=redact, # convenience callables appended after hooks
hooks_raise=True, # policy failures are fatal
)
让这个顺序刻意且记录在案,因为更晚的钩子会看到更早那个的改动。
限定作用域的埋点
只给需要它的代码挂上一个 tracer 或调试钩子,而不是重建整个 agent。hooks_installed 在退出时恢复
之前的列表,即便代码块抛异常也一样。
with agent.hooks_installed(DebugDump()):
agent.system_one(state, questions) # DebugDump only here
add_hook/remove_hook 不用代码块也能做同样的事,适合一个活得和进程一样久的 tracer。
进程级埋点
一个每个决策都该看到的 tracer 或指标钩子可以注册一次,而不是传给每个 Agent 和 Router。默认值
在实例钩子和逐调用钩子之前运行。
from laya import BaseHook, hooks
class Metrics(BaseHook):
def on_predict_end(self, ctx):
record(ctx.model, ctx.elapsed_ms)
hooks.set_default_hooks(hooks=[Metrics()])
这是全局状态,所以要有意识地限定它:在启动时设置一次,并在测试里 clear_default_hooks(),这样
一个测试就不会把钩子泄漏进下一个。
token 预算的塑造
一个 start 钩子可以为一次调用抬高 token 预算,例如当一个问题的选项很多、默认的 head 预算会让标签 塌掉时。有四个细节决定这个钩子是有帮助,还是悄悄把这次调用变得更糟:
- 一个 start 钩子的
ctx.head_max_len替换这次调用的预算。在它之前生效的是调用方自己的逐调用 值,或者ctx.agent.cfg里的 checkpoint 默认值 —— 所以要拿它做比较,而写一个光秃秃的数字可能 会拉低调用方已经设好的预算。 - 一次调用回答它携带的每一个问题,所以按其中最宽的那个定尺寸,而不是按碰巧最先出现的那个。
- 一旦选项装不下 head,
laya/common.py给每个选项max(4, (head_max_len - 16) // k)个 token。 因此16 + 4 * k正好落在那条下限上:每个标签仍然被削减到它和其他标签共享的那些 token,这 正是这个钩子写来避免的塌陷。16 + 8 * k让它们保持可区分。 - 状态拿到
max_len - head_max_len - 8个 token,所以拓宽的 head 必须连带拓宽max_len,否则 状态就丢了自己的窗口。
def widen_for_high_cardinality(ctx):
k = max((len(q.get("criteria", {}) or {}) for q in ctx.questions.values()), default=0)
if k < 50:
return
cfg = getattr(ctx.agent, "cfg", None) or {}
head = ctx.head_max_len if ctx.head_max_len is not None else cfg.get("head_max_len", 192)
window = ctx.max_len if ctx.max_len is not None else cfg.get("max_len", 512)
need = 16 + 8 * k # 8 tokens per label, not the core's floor of 4
if need > head: # only ever widen, never lower
ctx.head_max_len = need
ctx.max_len = max(window, need + 8 + 64) # 8 reserved, then room for the state
agent = laya.load("convaiinnovations/laya", on_predict_start=widen_for_high_cardinality)
这不碰共享的 agent 配置,所以并发的调用不受影响。同样的旋钮也可以逐调用使用:
agent.system_one(state, questions, head_max_len=512, max_len=1024)。
拓宽不是免费的:更长的窗口意味着更大的张量,而 checkpoint 是在 512(laya)和 1,024 个 token
上训练的。超过之后,用 predict_shortlist 收窄候选比硬撑预算更好。
反模式
阻塞的工作
钩子运行在调用线程上,而 laya.serve 用一个单独的推理 worker。一个 sleep、等一次网络往返或调用
input() 的钩子会把它后面的每个请求都堵住。
# bad: blocks the whole server
def audit(ctx):
requests.post("https://slow.example/decisions", json=..., timeout=30)
# better: enqueue, let a background worker ship it
def audit(ctx):
QUEUE.put_nowait(record(ctx))
如果你非要做慢工作,设置 hooks_concurrent=False 至少让钩子本身不重叠,并让 laya.serve 跑在
一个队列后面。
从 end 钩子抛出以控制流程
on_predict_end 在推理之后运行。在那里抛出会丢掉一个已经算出来的结果,而且在成功路径上还会浮现给
调用方。要拦截,用 start 钩子,在付推理代价之前;要改答案,改写 ctx.results。
没有锁的共享可变状态
同一个钩子实例在很多线程上运行。self.counter += 1 会竞争。
# bad
class Count:
def __init__(self): self.n = 0
def on_predict_end(self, ctx): self.n += 1
# good
import threading
class Count:
def __init__(self):
self.n = 0
self._lock = threading.Lock()
def on_predict_end(self, ctx):
with self._lock:
self.n += 1
静默失败
hooks_raise=False 每次失败警告一次,但一个自己捕获一切的钩子把真正的问题藏了起来。
# bad: no one will ever know the audit trail stopped
def audit(ctx):
try:
ship(record(ctx))
except Exception:
pass
如果一个钩子是可选的,让 hooks_raise=False 处理它,并盯着那些警告。如果不是,让它抛。
留住上下文
一个把 ctx 追加进列表的钩子会让整个状态、问题、结果和 agent 都活着。
# bad: unbounded memory growth
SEEN = []
def audit(ctx):
SEEN.append(ctx)
# good: keep only what you need
SEEN = []
def audit(ctx):
SEEN.append((ctx.run_id, ctx.model, ctx.elapsed_ms))
脱敏太晚
到了 on_predict_end,模型已经把状态 tokenize 过了。请在 on_predict_start 里脱敏。
在逐调用钩子里做逐问题逻辑
每次调用只有一个 PredictContext,一次前向传播回答每一个问题。没有逐问题的事件。在
on_predict_end 里遍历答案,也遍历状态:在一个批次上,一个上下文携带这次调用的每个状态。
def flag(ctx):
for result in ctx.results or []:
for qid, answer in result["answers"].items():
if answer.get("confidence", 1.0) < 0.5:
alert(qid, ctx.run_id)
递归 predict
一个调用 agent.predict/system_one 的钩子会再次运行钩子。没有深度保护的话它会递归下去。
# bad
def enrich(ctx):
ctx.results = [agent.predict(ctx.states[0], EXTRA_QUESTIONS)]
# good: guard, or use a separate agent with no hooks
def enrich(ctx):
if getattr(ctx, "_enriched", False):
return
ctx._enriched = True
ctx.results = [enricher.predict(state, EXTRA_QUESTIONS) for state in ctx.states]
hooks= 里的普通可调用对象
hooks= 接受钩子对象;一个裸的可调用对象说不清它是为什么事件,所以被拒绝。用
on_predict_start= / on_predict_end=。
# bad: TypeError
laya.load("convaiinnovations/laya", hooks=[lambda ctx: None])
# good
laya.load("convaiinnovations/laya", on_predict_end=lambda ctx: None)
在 end 钩子里假定结果存在
失败路径上 ctx.results 是 None,除非某个 start 钩子设过它。总要检查。
def audit(ctx):
if ctx.results is None:
log_failure(ctx.run_id, ctx.error)
return
log_success(ctx.run_id, ctx.results)
依赖顺序的钩子
一个从另一个钩子读取改动的钩子,除非顺序被钉住,否则很脆弱。已安装的钩子按列表顺序运行,然后是 便捷可调用对象;把任何耦合记录下来,或者把耦合的钩子合并成一个对象。