Skip to content

Latest commit

 

History

History
418 lines (300 loc) · 15.8 KB

File metadata and controls

418 lines (300 loc) · 15.8 KB

可观测性

MVP 必须

1. 白盒化回查

通过请求标识(如 uid / request_id)追踪某条请求的完整执行情况:

  • 经过了哪些算子
  • 每个算子的耗时
  • 每个算子的输入输出数据快照

引擎在每次 DAG 执行时自动记录这些信息。回查时按请求标识检索历史记录,还原完整的执行链路。

Trace 返回控制

引擎内部始终记录每个算子的 trace(名称、开始时间、耗时、是否 skip)。但 HTTP 响应中默认不返回 trace,以减少传输体积。

请求方通过 common 字段 _return_trace 控制:

{
  "common": {
    "_return_trace": true,
    ...
  },
  "items": []
}
  • _return_tracetrue → 响应包含 trace 数组
  • _return_trace 缺失或为 false(默认) → 响应不含 trace

_return_trace_ 开头,属于引擎保留字段,不会被算子读取,不需要在 common_input 中声明。

trace 数组仅包含实际执行或被 skip 的算子条目。当 DAG 因算子报错而中止时,未开始执行的下游算子不会出现在 trace 中。

2. 代码治理

自动检测和报告无用算子/分支:

  • Apple 侧(静态):配合 flow output 契约的死代码消除,在生成 JSON 时报错(已在 02_flow_abstraction.md 中定义)。
  • Pine 侧(运行时):统计每个算子和控制分支的实际执行情况,定期生成报告。长期未被执行的分支或算子标记为可清理候选。

3. 算子 debug 参数

所有算子都有一个可选的 debug 参数(默认 False),与 common_input / item_input / common_output / item_output 同级,属于通用参数。

开启后,该算子在运行时打印调试日志(输入数据、输出数据、耗时等详细信息),并在 trace 中填充 InputSnapshot / OutputSnapshot

全局 debug 开关

除逐算子 debug 外,还可在 Flow 级别全局开启:

flow = Flow("example", debug=True, ...)

编译到 JSON 根级 "debug": true。Go 侧也可通过 Option 覆盖:

engine, _ := pine.NewEngine(data, pine.WithDebug(true))

优先级:WithDebug Option > JSON debug 字段 > 逐算子 debug 参数。全局 debug 开启时,等效于对所有算子设置 debug=true

注意:debug 采集会复制 DataFrame 快照,对性能有影响,生产环境默认应关闭。

flow.transform_by_lua(
    common_input=["user_age"],
    item_input=["item_price"],
    function_for_item="adjust_price",
    item_output=["item_adjusted_score"],
    debug=True,  # 开启调试日志
    lua_script="""
        function adjust_price()
            return item_price * 0.8
        end
    """,
)

JSON 配置中体现为:

{
  "transform_by_lua_D4E5F6": {
    "type_name": "transform_by_lua",
    "$metadata": {
      "common_input": ["user_age"],
      "item_input": ["item_price"],
      "item_output": ["item_adjusted_score"]
    },
    "debug": true,
    ...
  }
}

Pine 执行时,debug: true 的算子输出详细日志,debug: false 或未设置的算子不输出。

引擎侧 debug 日志格式

引擎在 debug: true 的算子执行前后自动捕获输入/输出快照,并以 JSON 序列化格式打印到服务端日志。示例:

[pine:debug] operator="transform_by_lua_A1B2C3" duration=1.234ms
  input: {"common":{"user_age":16},"items":[{"item_price":100},{"item_price":200}]}
  output: {"item_writes":{"0":{"item_adjusted_score":80},"1":{"item_adjusted_score":160}}}

输入/输出均经过 json.Marshal 序列化,保证可读且可直接用于 diff/分析。序列化失败时回退到 %v

算子侧 debug 访问

Debug 是算子的配置属性(编译时确定),不属于请求输入数据。引擎通过 DebugAware 接口在编译阶段注入 debug 标志和算子名,算子通过嵌入 DebugHolder 获得 IsDebug()DebugLog() 能力。

MetadataAware / MetadataHolder 完全对称:

type MyOp struct {
    pine.MetadataHolder
    pine.DebugHolder     // 嵌入即自动实现 DebugAware
    threshold float64
}

func (o *MyOp) Execute(ctx context.Context, in *pine.OperatorInput, out *pine.OperatorOutput) error {
    if o.IsDebug() {
        o.DebugLog("custom state: threshold=%v, item_count=%d", o.threshold, in.ItemCount())
    }
    // ... 正常逻辑 ...
    return nil
}
  • o.IsDebug() bool — 返回当前算子的 debug 开关状态
  • o.DebugLog(format, args...) — 仅在 debug=true 时打印日志,自动附加算子名前缀 [pine:debug] operator="<name>"debug=false 时静默无操作

引擎在编译阶段调用 SetDebugInfo(operatorName, debug),算子无需自行构造 logger。

通用算子参数汇总

以下参数所有算子类型共有:

参数 必选 默认值 说明
common_input [] 读取的 common 字段
common_output [] 写入的 common 字段
item_input [] 读取的 item 字段
item_output [] 写入的 item 字段
common_defaults {} common 字段的缺失值默认值
item_defaults {} item 字段的缺失值默认值
debug False 是否打印调试日志

暂不计划

  • 全链路特征追踪:跨服务的数据血缘追踪。超出单引擎范畴,需分布式 tracing 基础设施支撑。

已实现

运行时指标体系

引擎通过 pkg/metrics 提供可插拔的指标接口(Counter / Gauge / Histogram),默认使用零开销的 Nop 实现。用户通过 pine.WithMetrics(provider) 注入自定义 Provider(如 Prometheus 适配器),即可将指标导出到外部监控系统。

同时,引擎内部保留 atomic 计数器作为独立的内置观测路径,通过 /stats JSON 端点始终可用,不依赖任何外部系统。

指标覆盖范围

调度器级

指标 类型 说明
pine_scheduler_runs_total Counter DAG 调度执行总次数
pine_operator_active Gauge 当前正在执行的算子数
pine_operator_exec_total Counter(operator) 算子成功执行次数
pine_operator_exec_duration_seconds Histogram(operator) 算子执行耗时分布
pine_operator_skip_total Counter(operator) 算子跳过次数
pine_operator_error_total Counter(operator) 算子失败次数

Lua pool 级

指标 类型 说明
pine_lua_pool_borrow_total Counter(operator) Lua state 借出总次数
pine_lua_pool_return_total Counter(operator) Lua state 归还总次数
pine_lua_pool_create_total Counter(operator) Lua state 创建总次数
pine_lua_pool_active Gauge(operator) 当前借出的 Lua state 数

配置热重载级

指标 类型 说明
pine_config_reload_total Counter 配置重载成功次数
pine_config_reload_errors_total Counter 配置重载失败次数
pine_config_reload_duration_seconds Histogram 配置重载耗时分布

资源级(redis_connection)

指标 类型 说明
pine_redis_pool_total_conns Gauge(name) 连接池中的连接总数(空闲 + 使用中)
pine_redis_pool_idle_conns Gauge(name) 连接池中的空闲连接数
pine_redis_ping_duration_seconds Histogram(name) 后台 PING 探针的往返耗时
pine_redis_up Gauge(name) 最近一次 PING 探针成功为 1,失败为 0

资源级指标由内置的 redis_connection 资源发出,仅当其 metrics_name 参数非空时才启用,标签 name 取该参数值;为空(默认)时不发出任何指标,也不会启动探针线程。探针间隔固定 15s,三运行时一致;资源初始化时立即执行一次探针,因此指标在第一个请求前即已就绪。

Metrics Provider 接口

import "github.com/Liam0205/pineapple/pine-go/pkg/metrics"

type Provider interface {
    NewCounter(opts MetricOpts) Counter
    NewGauge(opts MetricOpts) Gauge
    NewHistogram(opts HistogramOpts) Histogram
}

接口设计对齐 Prometheus mental model:With(labelValues...) 按位置传值、MetricOpts.LabelNames 声明标签名、Histogram.Observe(float64) 配合 metrics.DurationSeconds 转换。

Pineapple 核心库不依赖 prometheus/client_golang。Prometheus 适配器由用户在自己的项目中实现,约 80 行代码即可完成。

/stats 端点

GET /stats 返回复合结构:

{
  "operators": {
    "recall_static_A1B2C3": {"exec_count": 100, "skip_count": 0, ...},
    "transform_by_lua_D4E5F6": {"exec_count": 100, ...}
  },
  "scheduler": {"run_count": 100, "peak_concurrency": 4},
  "server": {"reload_count": 3, "reload_error_count": 0, "last_reload_duration_ns": 5234000},
  "operator_detail": {
    "transform_by_lua_D4E5F6": {"borrow_count": 100, "return_count": 100, "create_count": 8, "reuse_count": 93, "active_count": 0}
  },
  "resources": {
    "pine_redis_up": {"cache": 1},
    "pine_redis_pool_total_conns": {"cache": 2},
    "pine_redis_pool_idle_conns": {"cache": 1},
    "pine_redis_ping_duration_seconds": {"cache": {"count": 4, "sum_ns": 812000}}
  }
}
  • operators:per-operator 累计统计(exec/skip/error/duration)
  • scheduler:调度器级统计(运行次数、峰值并发)
  • server:配置热重载统计
  • operator_detail:实现 StatsProvider 接口的算子的自定义统计(如 Lua pool)
    • Lua pool 的 reuse_count 记录借出时命中池中既有 state 的次数;on-borrow miss = borrow_count - reuse_count,create_count 还额外计入构造时的一次预热创建。借此可区分「复用命中」与「新建」,评估池容量是否合适。该统计仅在 operator_detail 暴露,无对应 Prometheus 指标。
  • resources:资源级指标(metric-centric 形状,与 http 子树同构):{指标名: {标签值组合: 值}},counter/gauge 为标量、histogram 为 {count, sum_ns}(整数纳秒)。每层键按字典序排序,保证三运行时字节级一致。metrics_name 为空或无资源时该子树为 {}

资源级指标与 fan-out 路由

资源级指标走 fan-out(Tee) 路由,与引擎指标解耦。bundled server 给 ResourceManager 注入的不是裸 Provider,而是 Tee(注入的Provider, 专用Collector)

  • 每条资源指标同时写入调用方注入的 Provider(如 Prometheus 适配器)和专用 Collector
  • 引擎指标仍直接走注入的 Provider,不进入 Collector,因此 /stats.resources 子树只含资源级指标(当前为 4 个 redis 指标),不掺入 18 个引擎/服务级指标;
  • /stats.resources 读取 Collector 的聚合快照,无需外部 Prometheus 后端即可观测。

这样下游无需任何改动即可从 /stats.resources 直接读取新增指标;若已接入 Prometheus,同样的资源指标也会经注入的 Provider 导出,两条路径并存。Collector 与 Tee 随 ResourceManager 在热重载时一同原子替换,且生命周期长于它(资源持有指向 Tee 的裸指针)。

三运行时(Go pkg/metrics Collector+Tee、Java MetricsCollector+TeeProvider、C++ metrics::Collector+metrics::TeeProvider)行为字节级对齐,由 cross-validate section 16 锁定(正例 4 指标 + up=1 + ping count≥1,负例空 metrics_nameresources == {})。

Prometheus 接入示例

第三方项目实现 metrics.Provider 接口,约 80 行:

package promadapter

import (
    "github.com/prometheus/client_golang/prometheus"
    "github.com/Liam0205/pineapple/pine-go/pkg/metrics"
)

type provider struct{ r prometheus.Registerer }

func New(r prometheus.Registerer) metrics.Provider { return &provider{r: r} }

func (p *provider) NewCounter(opts metrics.MetricOpts) metrics.Counter {
    c := prometheus.NewCounterVec(prometheus.CounterOpts{
        Name: opts.Name, Help: opts.Help,
    }, opts.LabelNames)
    p.r.MustRegister(c)
    return &counter{vec: c}
}

// NewGauge, NewHistogram 类似...

type counter struct {
    vec *prometheus.CounterVec
    c   prometheus.Counter
}

func (c *counter) With(lvs ...string) metrics.Counter {
    return &counter{c: c.vec.WithLabelValues(lvs...)}
}

func (c *counter) Inc() {
    if c.c != nil { c.c.Inc() }
}

在 server wrapper 中注入:

mp := promadapter.New(prometheus.DefaultRegisterer)
server.Run(server.Config{
    ConfigPath: *configPath,
    Addr:       *addr,
    Metrics:    mp,
})

DAG 可视化

通过 Engine.RenderDAG(format) 方法或 GET /dag?format=dot|mermaid HTTP 端点获取 DAG 结构的可视化表示。

支持两种输出格式:

  • DOT (Graphviz):标准图描述语言,可通过 dot -Tsvg 渲染为 SVG/PNG。
  • Mermaid:可直接嵌入 Markdown / GitHub README,无需额外工具即可预览。

节点按算子类型着色,标签包含算子名和类型分类。

SubFlow 折叠渲染

当 pipeline 包含多个 SubFlow 且算子数量较多时,完整 DAG 可能难以阅读。引擎支持按层级将 SubFlow 折叠为聚合节点,只保留跨组的边(自动去重)。

API

// Level 0 = 全展开, Level 1 = 按顶层 SubFlow 折叠, Level 2 = 按两层折叠, ...
dot, _ := engine.RenderDAG("dot", pine.WithCollapse(1))
mmd, _ := engine.RenderDAG("mermaid", pine.WithCollapse(2))

HTTP

# 按顶层 SubFlow 折叠
curl http://localhost:8080/dag?format=dot&collapse=1

# 按两层折叠
curl http://localhost:8080/dag?format=mermaid&collapse=2

# 全展开(默认)
curl http://localhost:8080/dag?format=dot&collapse=0

折叠逻辑:

  • Node.SubFlow 路径(/ 分隔)的前 N 段分组,N 即 collapse 层级
  • 例如 collapse=1 时,recall/candidatesrecall 都归入 recall
  • collapse=2 时,recall/candidatesrecall 分属不同组
  • 无 SubFlow 归属的节点(SubFlow == "")保持独立
  • 跨组的边提升为聚合边,自动去重
  • collapse=0 或不传 collapse 参数时行为完全不变

引擎日志前缀

引擎支持为其诊断日志([pine-debug] 快照、observe_logLoggerAware 算子输出)统一添加前缀。前缀作用于引擎实例私有的 logger——多个引擎同进程时各自保留自己的前缀,进程全局 log 包不受影响(issue #172;此前是 log.SetPrefix() 全局生效、first-engine-wins 的旧语义)。

配置方式

Apple DSL 声明

flow = Flow(
    name="example",
    common_input=["user_id"],
    log_prefix="[my-service] ",
)

编译后 JSON 根级生成 "log_prefix": "[my-service] "

Go Option 覆盖

engine, err := pine.NewEngine(jsonConfig, pine.WithLogPrefix("[svc] "))

优先级:WithLogPrefix Option > JSON log_prefix 字段 > 无前缀。Option 是 nullable 三态(Go *string / Java nullable String / C++ std::optional):显式传空字符串也视为已设置,会覆盖 JSON 值。

log_prefix 存储在 Engine 实例上(Go Engine.Logger() / Java Engine.logPrefix() / C++ Engine::log_prefix()),作用域仅限该引擎的日志输出。

编程 API

engine, _ := pine.NewEngine(jsonConfig)
dot, _ := engine.RenderDAG("dot")       // Graphviz DOT
mmd, _ := engine.RenderDAG("mermaid")   // Mermaid flowchart

// SubFlow 层级折叠
dotCollapsed, _ := engine.RenderDAG("dot", pine.WithCollapse(1))
mmdCollapsed, _ := engine.RenderDAG("mermaid", pine.WithCollapse(2))

HTTP 端点

# DOT 格式(默认)
curl http://localhost:8080/dag

# Mermaid 格式
curl http://localhost:8080/dag?format=mermaid

# 渲染为 SVG(需要本地安装 graphviz)
curl -s http://localhost:8080/dag | dot -Tsvg -o dag.svg

# SubFlow 层级折叠
curl http://localhost:8080/dag?format=dot&collapse=1
curl http://localhost:8080/dag?format=mermaid&collapse=2