链路追踪
链路追踪
一次调用的每个钩子都收到同一个 PredictContext.run_id,所以一个 tracer 可以关联起始、结束和
错误事件(以及它打开的任何 span),而不用自己记账。
run_id
- 一个
uuid4().hex字符串,每次公开调用创建一次(predict_batch、system_one、Router.predict、ONNXAgent.system_one)。Router.predict_batch改为每个请求创建一个,即 该请求的一次Router.predict调用本会有的那个run_id。 - 被那次调用的每个钩子共享,包括
on_error和on_predict_end。 - 不是全局的,也不持久化:它在进程内标识一次调用。把它放进你的日志和出站载荷里,以便跨系统关联。
- 每次调用各不相同,所以两次调用永不冲突。
call A: run_id=1a2b... start ──┐
├─ end
error ─┘
call B: run_id=9f8e... start ─── end
一个最小 tracer
维护一个从 run_id 到你 start 时打开的 span 的映射,在 end 或 error 时关闭它。
import time
class Tracer:
def __init__(self):
self.spans = {}
def on_predict_start(self, ctx):
self.spans[ctx.run_id] = {
"model": ctx.model,
"started_at": ctx.started_at,
}
def on_predict_end(self, ctx):
span = self.spans.pop(ctx.run_id, None)
if span is None:
return
emit_span(
name="laya.predict",
run_id=ctx.run_id,
model=ctx.model,
duration_ms=ctx.elapsed_ms,
usage=ctx.usage,
ok=ctx.error is None,
)
def on_error(self, ctx):
# on_error runs before on_predict_end; leaving the span for on_predict_end is fine,
# or close it here if you prefer.
pass
agent = laya.load("convaiinnovations/laya", hooks=[Tracer()])
因为 on_predict_end 总会运行,它是关闭 span 的自然位置,而且它在失败路径上能看到 ctx.error。
span 的生命周期
on_predict_start ──► open span (run_id, model, started_at)
│
├─ inference
│
on_error ──► record ctx.error on the span
│
on_predict_end ──► close span (elapsed_ms, usage, ok)
错误 span
on_predict_end 在失败路径上运行,且 ctx.error 被设置,所以一个关闭位置就能处理两种情况:
def on_predict_end(self, ctx):
span = self.spans.pop(ctx.run_id, None)
if span is None:
return
if ctx.error is not None:
span["status"] = "error"
span["error_type"] = type(ctx.error).__name__
span["error_message"] = str(ctx.error)
span["duration_ms"] = ctx.elapsed_ms
emit(span)
如果你只实现 on_error,记住它在 on_predict_end 之前触发;不要两处都关闭 span,否则会重复计数。
OpenTelemetry
示例 examples/hooks/otel.py 记录计数器和一张直方图。要真正的
span,就从 tracer 里驱动 OTel API。钩子是同步的,所以用同步的 exporter(或者入队,再从一个 worker
导出)。
from opentelemetry import trace
tracer = trace.get_tracer("laya")
class OTelHooks:
def __init__(self):
self.spans = {}
def on_predict_start(self, ctx):
span = tracer.start_span("laya.predict", attributes={"laya.run_id": ctx.run_id, "laya.model": ctx.model})
self.spans[ctx.run_id] = span
def on_predict_end(self, ctx):
span = self.spans.pop(ctx.run_id, None)
if span is None:
return
if ctx.usage:
span.set_attribute("laya.input_tokens", ctx.usage["input_tokens"])
if ctx.error is not None:
span.record_exception(ctx.error)
span.set_status(trace.Status(trace.StatusCode.ERROR))
span.end()
laya.load("convaiinnovations/laya", hooks=[OTelHooks()], hooks_raise=False)
设置 hooks_raise=False,这样一个 tracer 出故障永远不会让请求失败。
分布式传播
run_id 是一个普通字符串,所以把它放进任何离开进程的东西里:日志行、发给审计服务的载荷,或者一次
决策触发下游调用时的 HTTP 头。
def audit(ctx):
requests.post(
"https://audit.example/decisions",
json=record(ctx),
headers={"X-Laya-Run-Id": ctx.run_id},
timeout=2,
)
记住像这样的阻塞调用会拖住调用线程;在并发服务时改成入队。见 模式。
嵌套调用
一个再次调用 predict 的钩子会启动一次新调用,带一个新的 run_id。父与子是独立的,除非你自己
把它们连起来。捕获父 id 并传下去:
def enrich(ctx):
for state in ctx.states: # a hook sees every state of the call
child = enricher.predict(state, EXTRA_QUESTIONS)
record_child_span(parent_run_id=ctx.run_id, child_run_id=child.get("run_id"))
一个父 run_id,每个状态一次子调用。一个读 ctx.states[0] 的写法只关联一个批次里的第一个决策,
并静默丢掉其余的。
防范递归(见反模式);最简单的防线是一个单独的、没有钩子的
enricher agent。