Skip to main content
可观测性不是一个独立进程,而是 Worker、Coordinator、Sync Invoker、Bundle 集群和 Python Executor 共享的一组 adapter 层基础设施。它由三个 Go 包组成:internal/obs(日志、指标、tracing)、internal/usage(计费用量采集)、internal/common/clientid(请求归属 ID 传播)。它们在系统中的位置见 架构总览

职责与边界

  • internal/obs 负责:安装进程级 slog handler、把 ctx 里的 task_id/call_id/instance_id 自动注入日志、声明所有 Prometheus 指标并提供 Record* 函数、初始化 OpenTelemetry tracer、启动 /metrics HTTP 端点。
  • internal/usage 负责:按 client_id 累计每个 task 收敛时的 executor 有效 CPU 时间,定期冻结成 record 写入 Kafka topic chaintable-usage
  • internal/common/clientid 负责:在 gRPC/HTTP 边界读写 client-id 与兼容 header x-instance-id,在进程内用 ctx 携带完整 ID。
  • 不负责:日志聚合栈、告警规则、Prometheus 服务端、OTLP collector。这些都在部署侧。
  • 不变量:core/ 包(Sans-IO)不 import internal/obs;只有 adapter 和 cmd/ 代码打日志、记指标、开 span。详见 设计原则
  • 不变量:日志 instance_id、chaintable-log trace_id、Prometheus instance_id label 都用 clientid.LogKey 去掉 instance: 前缀;usage 消息和所有下游 header 保留完整原值。
  • 不变量:Prometheus label 只允许闭合枚举(status/phase/backend/kind 等)加 instance_id/function_id 两个受 TTL 清扫器保护的高基数维度。不要用 task_id、call_id、URL、错误文本做 label。
  • 不变量:Metrics 用 promauto 注册到默认 registry;每个带 instance_idfunction_id label 的 vec 必须同时登记到 internal/obs/metrics_reaper.goinstanceIDVecs/functionIDVecs,否则 TestReaperVecListsMatchDeclarations 会失败。

代码位置

核心类型与接口

关键签名(从代码复制):

数据流

一次请求的归属标识从 gRPC metadata 进入,沿 ctx 流到日志、span、指标和 usage 四个出口:
1

边界归一化

internal/worker/adapters/grpc_server.goSubmitTaskinternal/syncinvoker/adapters/grpc_server.goInvoke/DebugInvokeinternal/coordinator/adapters/grpc_server.go 都以 ctx = clientid.WithContext(ctx, clientid.FromIncomingGRPC(ctx)) 开头。client-id 优先,缺失回退 x-instance-id,都缺失得到 unknown
2

进程内传播

Worker 在 task runtime 里用 clientid.WithContext(ctx, rt.instanceID) 重新绑定(internal/worker/adapters/orchestrator_dispatch.goorchestrator_phases.go),保证异步 goroutine 和 IO 回调不丢 ID。obs.WithInstanceID 只是 clientid.WithContext 的别名,不存第二份值。
3

日志与指标出口

ContextHandler.Handleclientid.LogKey 清洗后写 instance_idmetricInstanceID 对 Prometheus label 做同样清洗并刷新序列的 last-seen。chlog bridge 用 LogInstanceIDFromCtx 作为 trace_id
4

usage 与下游

task 收敛时 orchestrator_phases.goo.usageRecorder.Record(taskInstanceID, taskUsage),传完整 ID。internal/sdk/blockdblogicaltypesnoderpcmeta adapter 用 IntoOutgoingGRPC/IntoHTTP 双写两个 header 给下游。

日志

格式与字段

  • 每条日志一行 JSON,写到 stderrinternal/obs/logger.go)。time 是毫秒精度 UTC(2006-01-02T15:04:05.000Z07:00),与 Python executor 对齐,方便同一进程树的 Go/Python 日志混流解析。
  • 关联字段由 ctx 自动注入:task_idcall_idexecutor_idattempt_seq(绑定后即使为 0 也输出)、instance_id
  • 两条规则(来自 docs/specs/logging.md §4.1):有 ctx 的函数用 slog.InfoContext(ctx, ...);没有 ctx 的 goroutine 闭包必须接收调用方用 slog.Default().With(obs.AttrCallID, id) 预绑定的 *slog.Logger
  • 字段名一律走 obs.Attr* 常量。错误字段是 err
  • obs.Init() 不带参数;obs.BindProcess 没有任何 cmd/app 调用点,service / worker_addr 只出现在个别显式传参的日志里。

结构化日志类型(type 字段)

obs.LogType* 常量是 type 字段的全部合法值。用户视角最重要的几个: LogTypeTaskPhaseFinishedLogTypeCallFailedLogTypeIOError 常量仍在但不再单独输出,内容并入 task_finished

chaintable-log 桥接

设置 CHAINTABLE_LOG_BROKERS 后,obs.InitChlogContextHandler 的内层 handler 换成 chlogBridgeHandler:每条记录仍写 stderr,同时按 type → chlog WithType、其余 attrs → data payload、ctx 里的 instance id → trace_id 转发到 Kafka(默认 topic instance-logs)。没有 instance id 的记录(coordinator 后台日志等)被 SDK 静默丢弃。

指标

命名与注册

  • 命名空间 blockx,snake_case,counter 以 _total 结尾,直方图带单位后缀(_milliseconds_seconds_total_bytes)。
  • 声明方式统一为包级 promauto.New*Vec(...),通过 obs.Record* 函数暴露给 adapter;adapter 不直接持有 prometheus.* 对象。IO 自适应限流是例外,用自定义 prometheus.Collector 拉快照。
  • task 维度指标在 task 收敛时一次性观测:avg/p90/p99 由 internal/io/core/stats/ 的 CKMS 累积器在进程内算好,用 stat label 各 observe 一次。执行次数加权平均则用 *_duration_seconds_total / *_duration_samples_total 两个 counter 相除。

指标族与定义文件

internal/worker/adapters/call_metrics_test.goexecutor_metrics_test.gotask_barrier_metrics_test.go 是这些记录点的测试,不是定义文件。internal/worker/adapters/resource_metrics.go 是心跳用的 CPU/内存采样器(WorkerResourceSampler),不导出 Prometheus 指标。

最重要的指标

Tracing

  • obs.InitTracing(ctx, service) 始终设置 W3C TraceContext propagator。exporter 选择:BLOCKX_OTEL_SPAN_FILE 优先(JSON-lines 追加写文件,多进程共享)→ OTEL_EXPORTER_OTLP_ENDPOINT(OTLP/gRPC,insecure)→ 都没设则 noop。obs.TracingEnabled() 让热路径跳过 tracer.Start
  • 采样:默认 AlwaysSampleOTEL_TRACES_SAMPLER_ARG=0.1 这类值切到 ParentBased(TraceIDRatioBased)。生产必须显式设置。
  • 服务名由入口传入:blockx-workerblockx-bundle-workerblockx-coordinatorblockx-bundle-coordinatorblockx-syncinvokercmd/*/main.goProfile.Service),Python 侧是 blockx-executor
  • 自动 span:Coordinator→Worker gRPC 用 otelgrpc.NewClientHandler()/NewServerHandler()internal/coordinator/adapters/worker_rpc.gointernal/worker/app/app.go)。
  • 手工 span:Worker↔Executor UDS 用 internal/worker/adapters/executor/send.goudsTracer(tracer 名 github.com/Chaintable/blockx/worker/uds),span 名 uds.send.<MsgType>/uds.recv.<MsgType>/worker.io.dispatch;trace context 通过 MessageEnvelope.traceContext 字段过 UDS。Sync Invoker 用 obs.Tracer("blockx-syncinvoker")
  • BLOCKX_OTEL_UDS_EVENTS=1 打开 obs.AddUDSEvent/AddUDSEventAt 的细粒度 UDS pipeline event,仅 perf 调试用。
  • analyze_spans.py 读 span 文件,把 Go stdouttrace 与 Python ConsoleSpanExporter 两种 JSON 归一,按 trace 打印 span 树和 event 偏移。典型用法:
完整 span/event 清单见 blockx 仓库 docs/specs/tracing.md §4。

Usage 采集

  • 采什么:每个 task 收敛时 taskRuntime.snapshotUsage()internal/worker/adapters/task_runtime.go)把 executor 上报的有效执行时间(callEffectiveSum,毫秒)换成 usage.Sample{CPUTimeMicros}。不含 executor 排队、IO 等待、Builder/Result Handler 的 Go 侧 CPU。
  • 归属:Collector.Record(clientID, sample) 按完整 client_id 累计到 active;空 ID 记为 unknownCPUTimeMicros<=0 丢弃。
  • 往哪写:默认每 5s Flush,把 active 冻结成 []usage.Recordid UUID、client_idserviceresource_typeusage 毫秒、timestamp 冻结时刻),经 KafkaSink.Publish 同步写 topic chaintable-usage(key 为 client_idRequiredAcks=RequireAllMaxAttempts=1)。失败时 record 留在 pending,下轮以同一 id 重试;消费端按 id 去重。
  • 服务枚举:usage.ServiceBlockXWorker = "blockx_worker"cmd/worker)、usage.ServiceBlockXBundleWorker = "blockx_bundle_worker"cmd/bundle_worker),resource_type 固定 usage.ResourceTypeCompute = "compute"。由 Profile.UsageService 绑定,配置不能改。Sync Invoker 不创建 collector。
  • fail-closed:NewKafkaCollector 在未显式 Disabled 时要求 brokers 非空、topic 有分区,探测失败进程拒绝启动。

Client id 传播

状态与生命周期

  • usage CollectorStart 启动 ticker goroutine;Close 设置 closing barrier(之后 Record 被拒绝)、停 ticker、在调用方给的 shutdown ctx 内每 200ms 重试最终 Flush,然后 sink.CloseClose 返回错误视为计费关键错误,Worker 以非零退出。
  • Prometheus 序列 TTLinstance_id/function_id 的 label 值闲置超过 BLOCKX_METRICS_SERIES_TTL_MS(默认 24h)后由清扫器(默认每 10min,BLOCKX_METRICS_SERIES_SWEEP_INTERVAL_MS)用 DeletePartialMatch 删除,每轴每轮最多 256 个值。显式非正数关闭。清扫器随 StartPrometheusEndpoint 启动,即使没配抓取地址也会跑。
  • 表读写日志去重table_read 10 分钟窗口、table_write 1 分钟窗口,64 分片、每分片最多 2048 个 key。
  • tracing / UDS events 开关:进程启动时读一次环境变量,运行中修改无效。

监控端点

Worker 和 Coordinator 当前没有 HTTP 健康探针,也没有 gRPC health service。

配置

全部通过环境变量(Worker 也接受 JSON 配置文件的对应字段,见 internal/worker/app/config.goPrometheus/Chlog/Usage): 部署侧的完整变量清单见 blockx 仓库 docs/deploy.md §9,以及站内 部署概览

性能分析方法

  • 方法论:blockx 仓库 docs/specs/perf-methodology.md(建模 → 测量 → 比对 → 归因 → 决策的闭环);复现指引 docs/specs/perf-test-how-to-v2.md;结论在 perf-test-report-v2.md
  • 用例:e2e/perf/ 按 S0–S5 阶段前缀组织(s0_framework_*s1_builder_*s2_call_count_*s3_cpu_*s4_io_*s5_plugin_*),索引在 e2e/perf/README.md。产物落 e2e/perf/_artifacts/<TestName>/spans.jsonlpyspy/
  • 跑法:./test.sh -run TestPerf_XXX 在 cgroup(默认 4 vCPU / 8 GB)下跑;直接 go test ./e2e/perf/... 是无限制对照。带 with_trace 的子测试自动设置 BLOCKX_OTEL_SPAN_FILE
  • 线上压测:go run ./cmd/perf -coordinator host:port -n 100 -c 4 -mode full|meta|reserve,输出 P50/P90/P99。
  • 进程级归因:WORKER_CPU_PROFILE/WORKER_TRACE 写文件,或抓 /debug/pprof/profile
更多见 测试组织与命令

扩展点

  • 新增一个日志字段或 type:在 internal/obs/logger.goAttr*/LogType* 常量;如果是用户可见的 type,同步更新 docs/specs/blockx-log-types.md
  • 新增一个 Prometheus 指标:在 internal/obs/metrics.go(或对应子系统的 metrics_*.go)用 promauto 声明并加 Record* 函数;label 只能是闭合枚举加 instance_id/function_id;带这两个 label 的 vec 必须加入 metrics_reaper.goinstanceIDVecs/functionIDVecsinstance_id 值必须经过 metricInstanceID
  • 给新的 IO backend 加自适应指标:通过 internal/io/adaptive/observed 注册,它会调用 obs.RegisterIOAdaptiveMetrics(backend, provider)
  • 加 span:只在 adapter 层;先考虑能否用 event(obs.AddUDSEvent)表达;跨 goroutine 用 child span 而不是往已 End 的父 span 加 event。
  • 新增一个需要计费的入口:在 cmd/<entry>/main.goProfile.UsageService 传新的 usage.Service* 常量;service/resource_type 不允许从配置覆盖。
  • 新增一个下游 adapter:出站前调用 clientid.IntoOutgoingGRPC(ctx)clientid.IntoHTTP(req),不要自己读 metadata。

测试

  • internal/obs/handler_test.goContextHandler 的 slogtest 合规、字段注入、WithAttrs/WithGroup 保持注入、LogKey 清洗。
  • internal/obs/chlog_bridge_test.go:桥接双写、type 不进 data、group 前缀、trace_id 清洗。
  • internal/obs/metrics_test.gometrics_function_test.gostream_metrics_test.gometrics_syncinvoke_test.gometrics_io_adaptive_test.go:各 Record* 函数写出的 label 与数值。
  • internal/obs/metrics_reaper_test.gometrics_reaper_consistency_test.go:TTL 清扫器行为;TestReaperVecListsMatchDeclarations 解析源码,保证带 instance_id/function_id 的 vec 都登记到清扫列表。
  • internal/obs/table_dedupe_test.go:去重窗口与分片上限。
  • internal/usage/collector_test.gometrics_test.go:closing barrier、Kafka topic 探测、发送失败指标。
  • internal/common/clientid/client_id_test.go:header 优先级、双写、LogKey
  • 记录点测试在 internal/worker/adapters/call_metrics_test.goexecutor_metrics_test.gotask_barrier_metrics_test.go;进程级日志 schema 断言在 cmd/worker/worker_process_logging_test.go

相关文档

blockx 仓库 spec:
  • docs/specs/logging.md — 日志 schema、字段表、Go/Python 调用契约。
  • docs/specs/blockx-log-types.md — 用户视角的 type 日志清单、task_finished 字段定义、从日志推导指标的方法。
  • docs/specs/tracing.md — 启用方式、完整 span/event 清单、analyze_spans.py 用法、排障表。
  • docs/specs/client-id-propagation.md — header 优先级、出站双写、日志/监控清洗规则。
  • docs/specs/2026-07-13-blockx-usage-collection.md — usage 消息格式、CPU 口径、flush 与故障语义。
  • docs/specs/perf-methodology.mddocs/specs/perf-test-how-to-v2.md — 性能测试方法与复现。
  • docs/deploy.md §6、§9 — 部署侧可观测性配置与环境变量。
站内页面:WorkerCall 执行子系统IO 访问子系统Sync Invoker部署概览测试组织与命令