Skip to main content

BLOGE 0.9.8-RC1 在线版 · 事实校验 2026-09-15 · English

第 19 章 —— 生产环境中的可观测性

承诺: 读完本章后,你将知道如何在生产环境中看见一张 BLOGE graph 正在做什么 —— 通过 listener、interceptor、metrics、tracing、结构化日志、上下文传播与审计日志 —— 而不需要把这些逻辑写进 workflow 本身。


学习目标

  1. 解释两条主要的可观测性扩展点:负责生命周期事件的 ExecutionListener,以及负责包裹调用过程的 OperatorInterceptor
  2. 使用自定义 listener 捕获 graph / node 生命周期信号,并知道什么时候应该实现更细粒度的子接口。
  3. 正确地把 MetricsExecutionListenerTracingOperatorInterceptorLoggingExecutionListener 接到 GraphEngine 上。
  4. 使用 CallerContextCarrier 把 MDC 或 OpenTelemetry 上下文传播到引擎的虚拟线程中。
  5. 知道什么时候聚合指标就足够,什么时候必须使用 AuditJournalListener 提供的结构化逐事件视角。

前置条件

源示例

文件展示内容
GraphEngineListenerTest.java生命周期回调与 listener 异常隔离
ResilientOperatorWrapperObservabilityTest.javaretry、timeout 与 fallback 事件
MetricsExecutionListenerTest.javaMicrometer 定时器与计数器
TracingOperatorInterceptorTest.javagraph span 与 node span
AuditJournalListenerTest.java结构化审计条目与事件覆盖
BlogeObservabilityAutoConfiguration.javaSpring Boot 可观测性自动配置

为什么这很重要

一张 graph 今天在你的机器上跑得正确,仍然可能在生产环境里悄悄出问题:

  • 某个 node 在依赖升级后开始抛错。没有 metrics,你要等到用户投诉才知道。
  • 某个 retry 回路每次请求都烧掉 30 秒。没有逐节点耗时数据,你根本定位不到瓶颈。
  • 某条 分支改动让过去会执行的节点被跳过了。没有生命周期信号,系统看起来像是空闲,实际上却在做错误的选择。

BLOGE 将关注点明确分离:引擎 负责运行节点,扩展钩子 负责观察结果。正是这种分离,才让你能够在不把指标、trace 或日志代码塞进每个 operator 的前提下,获得生产可见性。


一次事故、一个 execution ID、四种投影

10:03,订单看板出现延迟尖峰。问题请求是 executionId=exec-order-42;它的 chargePayment 节点重试两次后进入 fallback。固定这一个身份,再切换问题:

图:一次事故的时间线

图:metric、trace、log 与 audit 四种投影

投影最适合回答的问题exec-order-42 的观察单独不能证明什么
Metric问题是否普遍?retry 数和 p95 延迟上升究竟哪次 attempt 失败
Trace时间花在哪里?chargePayment 子 span 占主要耗时持久事件的完整顺序
结构化 logruntime 报告了什么?timeout token 及 execution/node ID完整、有序的历史
Audit/event journal事件按什么顺序发生?start → retry → retry → fallback → complete未聚合时的全局发生率

关联键是主脊,不是第五种可观测产品。Metric 触发调查;同一个 execution ID 缩小 trace 与 log 范围;journal 复原顺序。这个事故是根据源码测试覆盖的事件 类型构造的教学复盘,不是从生产事故中截取的记录。


心智模型

把引擎想象成一个有固定聚光灯位点的舞台。每一次节点执行都会经过同一条可观测链路:

Diagram: 19-observability-in-production figure 1

一旦你知道自己需要哪个"聚光灯位点",就知道应该把可观测性能力接到哪里。


第一个可运行的示例

一个最小但生产味道很足的接线方式,是在普通引擎外面挂上 listener 和 interceptor:

var meterRegistry = new SimpleMeterRegistry();
var metrics = new MetricsExecutionListener(meterRegistry, "bloge");
var tracing = new TracingOperatorInterceptor(tracer);

var engine = GraphEngine.builder()
.registry(registry)
.listeners(List.of(metrics, tracing))
.interceptors(List.of(tracing))
.build();

这一处 builder 接线,就给了你同一次执行的三种视角:

  • graph / node 耗时指标,
  • graph 与 node 的 tracing span,
  • 驱动这两者的原始生命周期事件流。

关键不是具体用哪套库,而是契约:listener 观察事件,interceptor 包裹执行。


拆解分析

ExecutionListener —— 事件骨干

ExecutionListener 是一个复合接口,扩展了七个细粒度子接口:

子接口回调方法
GraphLifecycleListeneronGraphStartonGraphComplete
NodeLifecycleListeneronNodeStartonNodeCompleteonNodeFailedonNodeSkippedonNodeSuspendedonNodeResumed
ResilienceListeneronNodeRetryonNodeTimeoutonNodeFallback
StreamingListeneronNodeStreamStartonNodeStreamChunkonNodeStreamEndonNodeStreamError
IterationLifecycleListener循环/foreach 迭代事件
TimerEventListener定时器触发事件
EventCorrelationListener事件匹配事件

新的集成只需实现它所需的最窄子接口。已有代码可以继续实现聚合的 ExecutionListener

安全保证: 一个抛异常的 listener 不会破坏 graph 执行。引擎会捕获 listener 的异常,然后继续调用其余的 listener。这在 GraphEngineListenerTest.listener_exceptionInListener_doesNotBreakExecution 中有测试。

OperatorInterceptor —— AOP 层

OperatorInterceptor 包装每一次 operator 调用。拦截器接收一个 OperatorInvocation, 并必须调用 invocation.proceed() 来继续链路(或通过直接返回值来短路)。

public interface OperatorInterceptor {
Object intercept(OperatorInvocation invocation) throws Exception;
}

TracingOperatorInterceptor 就是这样为每个节点创建子 span 的 —— 它在 proceed() 之前启动 span,在 finally 块中结束。

CallerContextCarrier —— 跨虚拟线程的上下文传播

引擎在虚拟线程上运行每个节点。像 OpenTelemetry 上下文或 SLF4J MDC 这样的线程本地状态不会自动传播。 CallerContextCarrier 通过三阶段契约弥补了这一空白:

  1. capture() —— 在调用者线程上调用,快照上下文。
  2. Snapshot.attach() —— 在虚拟线程上调用,恢复上下文。
  3. Scope.close() —— 在节点完成时调用,清理上下文。 Scope.NOOP 可用于不需要清理的 carrier。

bloge-metrics-otel 中提供了两个实现:

生产可观测性组件

bloge-metrics-otel 模块提供了三个开箱即用的组件:

MetricsExecutionListener —— 发射 Micrometer 定时器和计数器:

  • bloge.graph.duration —— Timer,标签为 graph + outcome
  • bloge.node.duration —— Timer,标签为 graph + node + outcome
  • bloge.node.errors —— Counter,按异常类标记
  • bloge.node.retriesbloge.node.timeoutsbloge.node.fallbacksbloge.node.skipped —— Counter
  • bloge.stream.chunk.countbloge.stream.durationbloge.stream.errors —— 流式指标

所有指标名使用可配置的前缀(默认为 bloge)。

源码:MetricsExecutionListener.java

TracingOperatorInterceptor —— 在 graph 启动时创建一个父 span(bloge.graph.execute),并在每次 operator 调用时创建一个子 span(bloge.node.execute)。它 同时 实现了 OperatorInterceptorExecutionListener —— 将同一个实例注册到 interceptorslisteners 列表中。

源码:TracingOperatorInterceptor.java

LoggingExecutionListener —— 结构化 SLF4J 日志,使用 bloge.graphbloge.nodebloge.executionId MDC 键。两个可选标志 —— includeInputincludeOutput —— 控制载荷是否出现在日志中。除非你已审查过载荷中没有敏感数据,否则保持两者禁用。

源码:LoggingExecutionListener.java

审计日志

对于合规级别的事件记录,bloge-runtime-spi 提供了 AuditJournalListener。 它将 NODE_STARTNODE_COMPLETENODE_FAILEDNODE_RETRYNODE_SUSPENDNODE_RESUME 事件捕获为结构化的 AuditEntry 记录,并写入 AuditJournalStore。 异步刷新模式(虚拟线程 worker + 可配置批量大小)让热路径保持轻量。

AuditJournalListener 接受一个可选的 TimeSource 参数,使审计时间戳与引擎使用的时钟保持一致。在测试中你可以注入 ManualTimeSource 来获得确定性审计条目;在生产环境中默认使用 SystemTimeSource.INSTANCE

// 生产环境 —— 墙钟
var listener = new AuditJournalListener(store, config, jsonCodec);

// 测试 —— 确定性时间
var listener = new AuditJournalListener(store, config, jsonCodec, manualTimeSource);

AuditConfig 可控制:captureInputcaptureOutputasyncFlushflushInterval(默认 100ms)、batchSize(默认 64)。

ExecutionEventStore —— 租户级查询与清理

ExecutionEventStore 接口暴露了几个对生产可观测性管线至关重要的方法:

方法用途
loadEventsByTenant(tenantId, namespace, from, to, page, size)按租户/命名空间和时间范围分页查询事件
countEventsByTenant(tenantId, namespace)统计某个租户的事件总数 —— 适用于仪表盘指标
purgeEventsBefore(cutoff)删除早于截止时间点的事件
purgeEventsBefore(cutoff, limit)批量感知的清理,避免长时间的删除事务
purgeExecutionEvents(executionId)删除单次执行的所有事件(例如归档后)

所有新方法都提供了默认实现(抛出 UnsupportedOperationException),因此现有的 store 实现仍然可以编译。根据你的可观测性管线需求逐步升级即可。


按级别读取稳定日志标记

指标和 trace 只告诉你某个东西慢了。稳定 token 能定位具体 runtime 或 DSL fallback,又不会把周围散文变成合同。报警策略应依据实际级别与运行影响;FINE/FINER 只提供诊断线索,不应单独触发报警。

标记级别含义第一步排查
[TIMER_RESTORE_FAILED]SEVERETimerManager 无法重新加载 active durable timers;已持久化 timer 到下一次重启前不会触发。检查附带 exception、durable timer context、store 连接与被加载记录。
[STREAMING_ITEM_ERROR_SEND_FAILED]WARNINGstreaming foreach item 失败后,runtime 还因 output channel 拒绝发送而无法写出 error chunk。检查原始 item failure、channel 关闭竞态和 item-failure policy。
[SCHEMA_NUMBER_PARSE_FALLBACK]FINER数字强制转换在 integer 或 long 解析失败后尝试更宽类型。只用于 coercion 诊断;除非另有业务或错误率影响,不单独报警。
[ACCESSOR_METHOD_LOOKUP_MISS]FINER直接方法 accessor 不存在,DSL path access 正尝试 JavaBean getter fallback。若后续 path resolution 失败,再核对 property/getter 合同。
[SCHEMA_ENRICH_SKIPPED]FINEDSL schema enrichment 查询 Operator metadata 失败;编译会在 enrichment 减少的情况下继续。diagnostics 丢失 schema 细节时,检查附带 cause 与 registry metadata。
[FINGERPRINT_SKIPPED]FINEDSL compiler 无法计算 Operator fingerprint,并在缺少 fingerprint 时继续。在依赖恢复期 mismatch 检测前,检查 Operator 注册和 fingerprint 支持。

这些标记是稳定的,周围的文本不是。匹配方括号里的 token,不要匹配 散文。


执行事件日志(Execution Event Journal)

EventJournalListener 是一个内置的 ExecutionListener,把每一个引擎 级事件追加到一个只追加日志里。区别于审计日志(记录用于持久化的状态变更) 和 operator 拦截器链(只包单次节点调用),事件日志是你这次执行的 可回放记录。

ExecutionEventJournal journal = new InMemoryExecutionEventJournal();

GraphEngine engine = GraphEngine.builder()
.registry(registry)
.listeners(List.of(new EventJournalListener(journal)))
.build();

每条记录是一个 ExecutionEvent,带 execution id、单调序号、时间戳,以及 以下事件类型之一:

类型触发时机
GRAPH_STARTED / GRAPH_COMPLETED / GRAPH_FAILED引擎进入/离开一张图。
NODE_STARTED / NODE_COMPLETED / NODE_FAILED / NODE_SKIPPED单节点生命周期。
FOREACH_BATCH_STARTED / FOREACH_BATCH_COMPLETEDforeach 的每一批。
STREAM_ITEM_EMITTED / STREAM_COMPLETED流式 source。
WAIT_FOR_TIMER_ARMED / WAIT_FOR_TIMER_FIRED持久化 timer 生命周期。
LEASE_ACQUIRED / LEASE_REFRESHED / LEASE_LOST持久化租约生命周期。

Agent 事件 —— AgentEventJournalBridge

用上 bloge-agent-ext 模块时,AgentEventJournalBridge 把 agent 内部的 生命周期事件桥接到同一个日志里。这就给你一条统一的"图轮次 + agent 轮次 + 工具调用"回放时间线:

GraphEngine engine = GraphEngine.builder()
.registry(registry)
.listeners(List.of(
new EventJournalListener(journal),
new AgentEventJournalBridge(journal))) // 把 agent 事件也接入
.build();

bridge 增加的 agent 事件类型:

类型含义
AGENT_TURN_STARTED / AGENT_TURN_COMPLETEDagent 推理循环的一次迭代。
AGENT_LLM_CALLED发了一次 LLM 请求(带 token 用量和模型 id)。
AGENT_TOOL_INVOKED / AGENT_TOOL_COMPLETED / AGENT_TOOL_FAILED工具调用生命周期。
AGENT_MEMORY_TRIMMED某个 memory 策略丢弃或总结了较早的轮次。
AGENT_STREAM_CHUNK_EMITTEDstreaming agent 输出一个 chunk。

常见做法是:短期调试用内存 journal,长期回放挂一个 JdbcExecutionEventJournal 落到数据仓库。


常见陷阱

❌ 只将 TracingOperatorInterceptor 注册为 listener

var tracing = new TracingOperatorInterceptor(tracer);

var engine = GraphEngine.builder()
.registry(registry)
.listeners(List.of(tracing)) // graph span ✅
// missing .interceptors(List.of(tracing)) — no node spans! ❌
.build();

TracingOperatorInterceptor 有两个身份。ExecutionListener 这一侧负责在 onGraphStart / onGraphComplete 中创建 graph 级别 的 span;OperatorInterceptor 这一侧负责在 intercept() 中创建 node 级别 的子 span。若你只把它注册为 listener,你只会看到一条平坦的 graph span,而看不到逐节点的拆解。

修复方式:同一个实例同时注册到两个列表中。


引导式重写

打开 MetricsExecutionListenerTest.java,按下面三种方式改编它:

  1. 把指标前缀改为 myapp.bloge,并验证生成的 timer 名称确实跟着变化。
  2. 增加一条 fallback 路径,并断言 fallback 计数器正好递增一次。
  3. 注册一个内联 listener 来记录事件字符串,并把它看到的顺序与同一次执行发出的 metrics 对照起来。

这个练习的价值在于:graph 不变,但你不断切换观察它的镜头。


思维检查

  1. ExecutionListenerOperatorInterceptor 的区别是什么?

    (listener 响应生命周期事件;interceptor 直接包裹 operator 调用,并可在调用前后增加行为。)

  2. 为什么 TracingOperatorInterceptor 必须注册在两个地方?

    (它的 listener 一侧管理 graph 级 span;它的 interceptor 一侧创建 node 级 span。)

  3. CallerContextCarrier 解决了什么问题?

    (把调用者线程中的 MDC / OpenTelemetry 等上下文传播到引擎的虚拟线程里。)

  4. 什么情况下应该选择 AuditJournalListener 而不是普通 metrics?

    (当你需要结构化、逐事件、可审计的记录,而不是聚合计数器或定时器时。)

  5. 如果 listener 抛出异常,引擎提供什么保证?

    (listener 失败会被隔离;执行继续,其它 listener 仍然会运行。)


实验

目标: 用三种互补方式观察同一张 graph。

  1. 构建一张包含一个成功节点和一个不稳定节点的小 graph。
  2. 注册 MetricsExecutionListener,并验证出现了 retry 计数器。
  3. 在两个位置都注册 TracingOperatorInterceptor,并验证 graph 与 node span 都被创建。
  4. 使用内存版 store 注册 AuditJournalListener,并断言存储了预期的 NODE_STARTNODE_RETRYNODE_COMPLETE 事件。
  5. 通过 MdcContextCarrier 打开 MDC 传播,并验证节点日志能看到调用方的关联 ID。

进阶: 同一张 graph 分别在开启与关闭可观测性的情况下运行一次,然后比较每一种观察手段最擅长回答什么问题。


实验验收卡

  • 预期与观察: 同一 executionId 关联 metric、trace、log 和 audit。
  • 失败与恢复: 移除 correlation 字段;恢复统一 identity 后重建时间线。
  • 证明边界: 证明投影可关联,不证明指标解释业务原因。
  • 练习合同: 一次 retry→fallback 事故;只改 correlation;交付四投影时间线;全部回到同一执行即停止。

回顾

  • listener 观察生命周期事件;interceptor 包裹 operator 调用。
  • metrics、tracing、logging 与 audit journaling 都建立在同一组引擎扩展点之上。
  • CallerContextCarrier 是调用方上下文与虚拟线程执行之间的关键桥梁。
  • 最稳妥的可观测性集成是可复用的:它增加可见性,但不污染 graph 逻辑。

下一步

第 20 章——Spring 与生产环境接线中,你会把运行时可见性接进应用:条件化 bean、配置边界与面向运维的 endpoint。


参考链接

Coding Agent: Open the versioned task guide.