分析引擎与四大 DSL

这一篇解决什么
SkyWalking 有四种领域专用语言(DSL):OAL、MAL、LAL 管写入侧的指标聚合,MQE 管读取侧的查询计算。讲清楚每种 DSL 解决什么问题、语法长什么样、怎么用 ANTLR 解析 + Javassist 生成字节码(MQE 例外,用 Visitor 即时求值),以及 analyzer 模块群如何把 trace 切分成 Source。


全局视角:写入侧 vs 读取侧

OAP 的数据流是两条方向相反的流水线:

图1

OAL/MAL/LAL 把数据"算好"存进去,MQE 把存进去的数据"取出来"再加工。四者共享同一套"指标类型"和"存储模型",所以写入侧定义的指标名(如 service_resp_time)能直接在 MQE 里引用——这是它们能协同的前提,不是巧合。


OAL · Observability Analysis Language

解决什么问题

让用户用声明式语法声明"从哪类 Source 取哪个字段、用什么聚合函数聚合成什么指标",OAP 在启动时把这些声明编译成 Java 字节码并接入流式聚合流水线。

什么是"声明式聚合"
传统做法要为每个指标手写一个 Java 类(继承 Metrics)、手写一个 Dispatcher(把 Source 转成 Metrics)、手写 id()/hashCode()/serialize() 等样板方法,再注册到 MetricsStreamProcessor

OAL 让你只写一行 service_resp_time = from(Service.latency).longAvg();,编译器替你生成全部样板代码。这样 SkyWalking 内置的上百个指标都集中在 .oal 文件里,可读、可改、可关闭。

为什么要编译成字节码而不是解释执行?流式聚合是高频热路径,每秒可能处理几十万条样本,解释执行的开销扛不住;编译成 Java 类后跑的是原生方法调用,和手写代码零差别。

真实脚本示例

oap-server/server-starter/src/main/resources/oal/core.oal(核心 OAL,所有 service/endpoint/关系拓扑指标都在这里)节选:

// oal/core.oal:20
service_resp_time = from(Service.latency).longAvg().decorator("ServiceDecorator");
// oal/core.oal:21
service_sla = from(Service.*).percent(status == true).decorator("ServiceDecorator");
// oal/core.oal:23
service_percentile = from(Service.latency).percentile2(10); // p50/p75/p90/p95/p99
// oal/core.oal:29
service_relation_client_cpm = from(ServiceRelation.*).filter(detectPoint == DetectPoint.CLIENT).cpm();
// oal/core.oal:58
endpoint_mq_consume_latency = from((str->long)Endpoint.tag["transmission.latency"])
                                .filter(type == RequestType.MQ)
                                .filter(tag["transmission.latency"] != null).longAvg();

语法结构(见 OALParser.g4):

  • 变量名 = from(Source.字段).filter(条件).聚合函数().decorator("装饰器");
  • Source 是固定枚举(ServiceEndpointServiceRelationServiceInstanceDatabaseAccessCacheAccessMQAccess…,见 OALParser.g4:56source 规则)。为什么是固定枚举而不是任意类型?因为 Source 是 trace 切分的产物,字段在编译期必须确定,enricher 要靠反射查 Source 类字段来生成 Dispatcher 代码(见下文)。
  • 聚合函数:longAvg()cpm()percent(status==true)percentile2(10)sum()count()apdex(name,status) 等,每个对应 server-core 里一个度量基类(如 LongAvgMetrics)。
  • (str->long) 是类型转换(castStmtOALParser.g4:214),把字符串 tag 转成 long 再算平均。tag 来自探针上报的 key-value,原始都是字符串,算数值聚合前必须显式转换。
  • disable(segment); 可关闭某条内置流(见 oal/disable.oalOALParser.g4:36disableStatement)。

ANTLR 语法与 parser 生成

  • 语法定义:oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALLexer.g4OALParser.g4

ANTLR / AST 速成
ANTLR 是"语法规则 → 解析器代码"的生成器。你写 .g4 描述语法(lexer 词法 + parser 句法),构建时自动生成 OALLexer.java/OALParser.java 等类。生成的 parser 把一行 OAL 文本解析成"解析树"(parse tree)。

AST(抽象语法树)是把解析树里"对代码生成有用的信息"提炼成精简的 Java 对象。OAL 里这个 AST 就是 MetricDefinition(不可变,字段全 final,见 oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/model/MetricDefinition.java:39),由 OALListenerV2oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/parser/OALListenerV2.java:68)在 exitAggregationStatement:123)里用 Builder 拼出。

编译与执行原理

图2

入口在 server-core,引擎实现在 oal-rt。这条链路分七步:

  1. 触发OALEngineLoaderService.load(OALDefine)oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/oal/rt/OALEngineLoaderService.java:58)。各 receiver 在 start() 里调用它,传入含 .oal 路径 + Source 包名的 OALDefine。核心 OAL 由 CoreOALDefine.INSTANCE 触发(指向 oal/core.oal)。每个 OALDefine 只激活一次(:60oalDefineSet.contains 去重)。
  2. 反射加载引擎:因为 server-core 在 Maven 反应堆里先于 oal-rt 编译,server-core 不能直接 import oal-rt 的类,只能用反射创建 OALEngineV2OALEngineLoaderService.java:88Class.forName + getConstructor + newInstance)。不这么做会形成 Maven 模块循环依赖,编译不过。
  3. 解析OALEngineV2.start()oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/OALEngineV2.java:90)读 .oalOALScriptParserV2.parse()oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/parser/OALScriptParserV2.java:99):建 OALLexer:101)→ CommonTokenStream:104)→ OALParser:107)→ parser.root():115),再由 ParseTreeWalkerOALListenerV2OALScriptParserV2.java:123)把解析树转成 List<MetricDefinition>
  4. enrichMetricDefinitionEnricher.enrich()oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/generator/MetricDefinitionEnricher.java:96)通过反射查 Source 类字段,给每个 metric 补上"要复制哪些字段、哪些是持久化字段、scopeId"等元信息,输出 CodeGenModel。这一步把 AST 从"语法描述"升级成"代码生成所需的全部上下文"。
  5. 代码生成OALClassGeneratorV2.generateClassAtRuntime()oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/generator/OALClassGeneratorV2.java:143)用 Javassist(运行时字节码操作库)+ FreeMarker 模板/code-templates-v2/*.ftl)为每条 metric 生成三类 Java 类:
    • Metrics 类generateMetricsClass:196):继承对应聚合基类(如 LongAvgMetrics),加 @Column/@Stream 注��、id()/serialize()/deserialize()/toHour()/toDay() 方法。类名如 ServiceRespTimeMetrics
    • MetricsBuilder 类generateMetricsBuilderClass:341):实现 StorageBuilder,做实体↔存储互转。
    • Dispatcher 类generateDispatcherClass:393):实现 SourceDispatcher<Service>,把 Source 对象分发到对应 Metrics。
  6. 接线OALEngineV2.notifyAllListeners()��OALEngineV2.java:134)把生成的 Metrics 类交给 StreamAnnotationListener,Dispatcher 类交给 SourceReceiverDispatcherDetectorListener
  7. 接入流式聚合StreamAnnotationListener.notifyoap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/StreamAnnotationListener.java:50)读生成类上的 @Stream 注解,按 processor 路由到对应 StreamProcessor:58-60 分支 Record/Metrics)。OAL 生成的 metric 都标了 MetricsStreamProcessor.class,于是进入 MetricsStreamProcessor.createoap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java:165,kind 默认 OAL),构建 worker 链。

代码生成的小贴士
这些类是运行时在内存里拼出 .class 字节码加载进 JVM 的,磁盘上看不到 .java。设环境变量 SW_DYNAMIC_CLASS_ENGINE_DEBUG=true 会把 .class + .java 旁车文件写到 oal-rt/ 目录便于调试(见 oal-rt/CLAUDE.md 的 “Debug Output”)。


MAL · Meter Analysis Language

解决什么问题

让用户用类 PromQL/Groovy 风格的表达式把"接收到的原始 meter 样本(Prometheus/OTel metric)"加工成 SkyWalking 的指标。

OAL vs MAL 的区别

OAL MAL
输入 OAP 内部 Source(trace 切分产物) 外部上报的 metric 样本(SampleFamily,带标签)
语法 from(Source.field).func() 声明式 metric.sum(['tag']).rate('PT1M').service(['host'], Layer.OS_LINUX) 表达式链 + 闭包
配置载体 .oal 纯文本 YAML(metricPrefix/expSuffix/filter/metricsRules),exp 字段是 MAL 表达式
触发 OALEngineLoaderService.load(OALDefine) receiver(如 OTel)启动时 Rules.loadRules + new MetricConvert(rule, meterSystem)
落盘注册 StreamAnnotationListenerMetricsStreamProcessor MeterSystem.createMetricsStreamProcessor(同一 processor,kind=MAL)

OAL 的 Source 是固定枚举(Service/Endpoint…),但外部系统(VM、MySQL、Redis、K8s…)上报的 metric 千变万化,OAL 无法表达;MAL 提供一套"对样本做标签过滤、四则运算、聚合、直方图/分位数"的表达式语言来填补。

真实脚本示例

MAL 表达式嵌在 YAML 里。oap-server/server-starter/src/main/resources/otel-rules/vm.yaml(VM node 监控)节选:

# otel-rules/vm.yaml:31-38
filter: "{ tags -> tags.job_name == 'vm-monitoring' }"   # 闭包:只处理 job_name=vm-monitoring 的样本
expSuffix: service(['node_identifier_host_name'], Layer.OS_LINUX)  # 公共后缀:绑定到 service scope
metricPrefix: meter_vm                                   # 指标名前缀
metricsRules:
  - name: cpu_total_percentage
    exp: (node_cpu_seconds_total * 100).tagNotEqual('mode','idle').sum(['node_identifier_host_name']).rate('PT1M')
  - name: memory_used
    exp: node_memory_MemTotal_bytes - node_memory_MemAvailable_bytes

最终指标名 = metricPrefix + "_" + name = meter_vm_cpu_total_percentage;最终表达式 = expPrefix + exp + expSuffix(见 MetricConvert.formatExpoap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/MetricConvert.java:306)。把 expSuffix 抽到文件级,是为了让同文件的多条规则共享同一 scope 绑定,避免每条重复写 service(['host'], Layer.OS_LINUX)

MAL 表达式语法(oap-server/analyzer/meter-analyzer/src/main/antlr4/org/apache/skywalking/mal/rt/grammar/MALParser.g4)支持:

  • 算术:metric1 + metric2(metric * 100)additiveExpression:47/multiplicativeExpression:51)。
  • 方法链:metric.sum(['tag']).rate('PT1M').histogram().histogram_percentile([50,75,90,95,99])postfixExpression:64)。
  • 闭包(Groovy 风格):{ tags -> tags.job_name == 'vm-monitoring' },用于 tag({...})forEach([...], {...})、文件级 filterclosureExpression:126)。为什么用 Groovy 风格闭包而不是 Java lambda?因为 Javassist 不支持 lambda 语法,闭包要单独生成"伴生类"(MalClosureCodegen),Groovy 风格的 { args -> body } 更容易被 Javassist 拼出。
  • 扩展函数 SPI:metric.genai::estimateCost()命名空间::方法():79IDENTIFIER DOUBLE_COLON IDENTIFIER),通过 java.util.ServiceLoader 发现,方便第三方插件加自定义函数。

编译与执行原理

图3

核心类(均在 oap-server/analyzer/meter-analyzer/.../v2/):

  1. 加载 YAMLRules.loadRules("otel-rules", enabledRules)oap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/prometheus/rule/Rules.java:60)用 SnakeYAML(:149new Yaml().loadAs)把每个 YAML 解析成 RuleOpenTelemetryMetricRequestProcessor.startoap-server/server-receiver-plugin/otel-receiver-plugin/.../otlp/OpenTelemetryMetricRequestProcessor.java:218)在 OTel receiver 启动时调它(:226),对每个 rule new MetricConvert(rule, meterSystem):237)。
  2. 每条规则建一个 Analyzer,分两阶段MetricConvert.java:107 构造函数):
    • Phase 1 prepareAnalyzer.prepareoap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/Analyzer.java:176)调 DSL.parse():197)→ MALScriptParser.parseoap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/compiler/MALScriptParser.java:213,内部建 MALLexer:215MALParser:217parser.expression():233,由 MALExprVisitor:280)转 AST)得到 MALExpressionModel.Expr,再 MALClassGenerator.compile.../compiler/MALClassGenerator.java:150)用 Javassist 生成实现 MalExpression 接口的类。闭包另生成"伴生类"(MalClosureCodegen)。同时 MALMetadataExtractor.extractMetadata.../compiler/MALMetadataExtractor.java:55)静态从 AST 提取 ExpressionMetadata(输入 metric 名、scope、标签),供运行时 O(1) 过滤。
    • Phase 2 registerAnalyzer.registerAnalyzer.java:218)调 MeterSystem.create(metricName, functionName, scopeType, ...)oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/MeterSystem.java:113),再用 Javassist 生成存储侧 Metrics 子类并接入 MetricsStreamProcessor.create(..., MetricStreamKind.MAL)MetricsStreamProcessor.java:174)——和 OAL 走的是同一个 MetricsStreamProcessor,只是 kind 标记为 MAL
    • 为什么要分两阶段?Phase 1 只编译不落盘,Phase 2 才触发存储 DDL。这样同文件某条规则编译失败时,前面已编译的规则不会留下"已建表但无数据"的脏状态,便于回滚(见 MetricConvert.java:166-183registered 跟踪集)。
  3. 运行时执行:OTel receiver 收到一批 metric(按 Prometheus metric 名组织成 ImmutableMap<String, SampleFamily>),调 MetricConvert.toMeterMetricConvert.java:325)广播给每个 Analyzer.analyseAnalyzer.java:264):先用编译期提取的输入 metric 名做 O(1) 过滤(:266-270,避免对每条规则全量遍历),再过文件级 filter 闭包(:297),再调生成的 expression.run(input):325)算出 Result,最后按 metric ��型(single/labeled/histogram/histogramPercentile)调 meterSystem.buildMetrics + doStreamingCalculation 落盘。

关键点
MAL 和 OAL 最终都汇入同一个 MetricsStreamProcessor,差别只在"Source 怎么来"和"Metrics 类怎么生成"——OAL 由 OAL 引擎从 .oal 生成,MAL 由 MeterSystem 从 MAL 表达式生成。下游 worker 链、存储、查询对两者完全一致,这�� SkyWalking 能统一内部 trace 指标和外部 meter 指标的关键。


LAL · Log Analysis Language

解决什么问题

让用户用 DSL 声明"如何解析一条日志正文(JSON/YAML/正则)、从解析结果里提取 service/endpoint/layer/timestamp/tags、可选地从日志里提取指标(metrics 块)、以及如何 sink(保存/丢弃/采样)“。LAL 把"无结构的日志流"变成"有结构的可查询日志 + 从日志里抽出来的指标”。

真实脚本示例

oap-server/server-starter/src/main/resources/lal/nginx.yaml

# lal/nginx.yaml:16-60
rules:
  - name: nginx-error-log
    layer: NGINX
    dsl: |
      filter {
        if (tag("LOG_KIND") == "NGINX_ERROR_LOG") {       # 条件分发
          text {                                          # 文本解析器
            regexp $/(?<time>\\\\d{4}/\\\\d{2}/\\\\d{2} \\\\d{2}:\\\\d{2}:\\\\d{2}) \\\\[(?<level>.+)].*/$
          }
          extractor {                                    # 字段提取
            tag level: parsed.level
            timestamp parsed.time as String, "yyyy/MM/dd HH:mm:ss"
            metrics {                                     # 从日志提取指标
              timestamp log.timestamp as Long
              labels level: parsed.level, service: log.service, service_instance_id: log.serviceInstance
              name "nginx_error_log_count"
              value 1
            }
          }
          sink { }                                        # 保存(空 sink = 全保存)
        }
      }

最简单的"保存所有日志"(lal/default.yaml):

rules:
  - name: default
    layer: GENERAL
    dsl: |
      filter {
        sink {
        }
      }

LAL 语法结构(oap-server/analyzer/log-analyzer/src/main/antlr4/org/apache/skywalking/lal/rt/grammar/LALParser.g4):

  • 顶层 root: filterBlock EOF:38),filterBlock: FILTER L_BRACE filterContent R_BRACE:43),filterContent: filterStatement*:47)。
  • 三类块:parserBlocktext{}/json{}/yaml{}:62)、extractorBlock:111)、sinkBlock:227)。
  • extractor 里可设 service/instance/endpoint/layer/traceId/timestamp/tag/metrics/def(局部变量)。
  • metrics 块(:186metricsBlock)声明 name/timestamp/labels/value——这就是 LAL→MAL 的桥接点
  • sink 里可放 sampler { rateLimit("id") { rpm 6000 } }(限流采样)、enforcer{}(强制保存)、dropper{}(丢弃)(sinkStatement:235,分 samplerBlock:242/enforcerStatement:262/dropperStatement)。
  • if (condition) { ... } else if ... else ... 控制流,条件支持 tag("KEY") 函数、parsed?.field 安全导航(?.)。

编译与执行原理

核心类(oap-server/analyzer/log-analyzer/.../v2/):

  1. 加载与启动LogAnalyzerModuleProvider.startoap-server/analyzer/log-analyzer/.../provider/LogAnalyzerModuleProvider.java:154)调 factory.loadStaticRules(),把 lal/ 目录下每个 YAML 的每条 rule 编译。同时把 log-mal-rules/ 下的 YAML(从日志抽指标的 MAL 规则)new MetricConvert(rule, meterSystem):160)注册进 converter map——所以 LAL 模块同时承载 lallog-mal-rules 两个 catalog。
  2. 解析 DSLLALScriptParser.parseoap-server/analyzer/log-analyzer/.../compiler/LALScriptParser.java:86)用 ANTLR 解析:建 LALLexer:87)→ CommonTokenStream:88)→ LALParser:89)→ parser.root():105),转成 LALScriptModel(不可变 AST)。LALConfig 是单条 rule 的配置 POJO(name/dsl/layer/inputType/outputType)。
  3. 代码生成DSL.compileoap-server/analyzer/log-analyzer/.../dsl/DSL.java:107)调 LALClassGenerator.compile 为每条 LAL rule 生成单个类(实现 LalExpression,方法 execute(FilterSpec, ExecutionContext)),extractorsink 块编成私有方法 _extractor()/_sink()。编译期还会做"数据源类型分析":根据有无 json{}/yaml{}/text{} 决定 parsed.* 生成哪种 getter。为什么是单类而不是像 OAL 那样三类?LAL 的输入是单条日志,没有 Source→Metrics 的分发关系,一个 execute 方法顺序跑 parser→extractor→sink 就够了。
  4. 运行时执行:日志 receiver 收到 LogData,对每个 layer 用 LogFilterListener.Factory.create(layer)oap-server/analyzer/log-analyzer/.../provider/log/listener/LogFilterListener.java:539)建一个 LogFilterListenerLogFilterListener.java:66)。两阶段:parse():114)建 ExecutionContextbuild():89)对每个 DSL 调 dsl.evaluate(ctx)DSL.java:134)→ 生成的 LalExpression.execute()DSL.java:142)→ 依次跑 filterSpec.json/text/yaml(ctx)_extractor()filterSpec.sink(ctx)ExecutionContext 显式传参,无 ThreadLocal——这样能在多线程并发处理日志时避免线程间状态污染。

LAL→MAL 桥
extractor 里的 metrics {} 块由 MetricExtractoroap-server/analyzer/log-analyzer/.../dsl/spec/extractor/MetricExtractor.java)处理:prepareMetrics:67)建一个 SamplesubmitMetrics:74)把 Sample 包成 SampleFamily:78-79),然后调 provider.getMetricConverts() 拿到 log-mal-rules 的所有 MetricConvert,调 it.toMeter(map):97)——这一步把 LAL 抽出来的样本直接喂给 MAL pipeline,复用 MAL 的 Analyzer 算出指标。

所以"从日志里提取指标"是 LAL(解析+提取样本)+ MAL(聚合样本成指标)的接力。


MQE · Metrics Query Engine

解决什么问题

让用户在查询时用一条表达式把"已经聚合好、存进存储的指标"读出来再做二次计算(四则运算、聚合、趋势 increase/rate、TopN、排序、重命名标签、布尔逻辑、基线对比)。MQE 是读取侧的表达式语言,对应写入侧的 OAL/MAL。

与 OAL 的关系
OAL/MAL 写入,MQE 读出。OAL/MAL 定义并聚合了 service_resp_time 这样的指标,MQE 让你写 service_resp_time / service_cpmincrease(http_count[5m]) 这样的表达式,OAP 读出底层指标再做运算返回给 UI/告警。

真实脚本示例

MQE 没有 .oal 那样的"配置文件",它是查询时由用户(GraphQL metricsExpression 查询、告警规则)传入的字符串。语法见 oap-server/mqe-grammar/src/main/antlr4/org/apache/skywalking/mqe/rt/grammar/MQEParser.g4。典型表达式:

service_resp_time + service_cpm                    # 两个指标相加(addSubOp,:29)
relabels(metric, label='a', replaceLabel='b')      # 重命名标签(relablesOP,:39)
aggregate_labels(sum, metric, [tag])               # 跨标签聚合(aggregateLabelsOp,:40)
increase(metric, 5)                                 # 趋势计算(trendOP,:35)
top_n(metric, 10, des)                              # TopN(topN,:81)
sort_values(metric, 10, des, avg)                  # 排序(sortValuesOP,:41)
service_resp_time > 100 ? service_cpm : 0          # 布尔/比较

支持:标量 scalar:66)、metric(可带 labelList 过滤,:56)、聚合 avg/count/latest/sum/max/min、数学函数 abs/ceil/floor/round、趋势 increase/rate、逻辑 view_as_seq/is_present、布尔 and/or���比较 ==/!=/>/...</>=topNOf 合并多个 TopN、baseline(AI 预测基线,需 ai-pipeline 模块)。

编译与执行原理

MQE 与 OAL/MAL/LAL 的根本区别
MQE 不做代码生成,而是用 ANTLR 的 Visitor 模式直接在内存里"边解析边执行"。因为查询是高频、短生命周期、且要支持错误返回,没必要编译成字节码——编译一次的开销远超一次查询的执行时间,缓存这些临时类还会撑爆 PermGen。

图4

  1. 入口MQEExecutor.execute(expression, entity, duration, foreign)oap-server/server-query-plugin/query-graphql-plugin/.../mqe/rt/MQEExecutor.java:71)。GraphQL 查询 MetricsExpressionQuery 调它;告警侧 AlarmMQEVisitor 也复用同一套语法。
  2. 解析:建 MQELexer:83)→ MQEParser:85)→ parser.expression():89)得到 parse tree,错误由 ParseErrorListener 捕获(:84/:86)。
  3. Visitor 执行MQEVisitoroap-server/server-query-plugin/query-graphql-plugin/.../mqe/rt/MQEVisitor.java:65)继承 MQEVisitorBaseoap-server/mqe-rt/src/main/java/org/apache/skywalking/mqe/rt/MQEVisitorBase.java:64,一个 MQEParserBaseVisitor<ExpressionResult>)。每个语法备选对应一个 visit* 方法:visitAddSubOpMQEVisitorBase.java:110)→ BinaryOp.doBinaryOp:124),visitAggregationOp:178)→ AggregationOp.doAggregationOp:188),visitTrendOP:404)→ TrendOp.doTrendOp:423)。递归 visit(ctx.expression()) 自底向上算,最终返回 ExpressionResult。把运算逻辑拆到 BinaryOp/AggregationOp/TrendOp 这些独立类,是为了让告警侧的 AlarmMQEVisitor 能复用同一套运算实现,不重复造轮子。
  4. 叶子节点 visitMetric:这是 MQE 真正"读存储"的地方(MQEVisitor.java:110)。MQEVisitor 持有 MetricsQueryService:69,通过 ModuleManager 拿到 CoreModule 查询服务),叶子 metric 节点调 getMetricsQueryService().readMetricsValues:280)/readLabeledMetricsValues:311)把底层存储的指标读出来,包装成 ExpressionResult 往上传。
  5. 调试:整个执行在 DebuggingTraceContext 里跑,便于 UI 展示表达式求值过程。

对比 OAL
OAL 是"声明式 → 编译成 Java 类 → 长期运行聚合 worker";MQE 是"表达式 → ANTLR 解析 → Visitor 即时求值 → 返回结果",无落盘、无代码生成。


四种 DSL 总览对比

DSL 方向 解决问题 语法载体 编译方式 运行时落点
OAL 写入(聚合) 声明式从内部 Source 聚合指标 .oal 文本 ANTLR→AST→Javassist+FreeMarker 生成 Metrics/Builder/Dispatcher 类 MetricsStreamProcessor(经 @Stream 注解 + StreamAnnotationListener
MAL 写入(聚合) 表达式处理外部 meter 样本 YAML(exp 字段) ANTLR→AST→Javassist 生成 MalExpression 类 + 闭包伴生类 MeterSystem.createMetricsStreamProcessor(kind=MAL)
LAL 写入(解析+提取) 解析日志、提取字段、从日志抽指标 YAML(dsl 字段) ANTLR→AST→Javassist 生成 LalExpression 类(单类+私有方法) LogFilterListenerexecute()metrics{} 块经 MetricExtractor 喂给 MAL
MQE 读取(查询) 查询时对已存指标做表达式计算 查询字符串(无配置文件) ANTLR→Visitor 直接执行(无代码生成) MQEVisitor 递归 visit,叶子 visitMetricMetricsQueryService 读存储

链路:探针 → agent-analyzer 切分成 Source / OTel → MAL 样本 / 日志 → LAL 解析 → OAL 或 MAL 聚合成指标 → MetricsStreamProcessor 落盘 → 用户写 MQE → MQEVisitor 读存储 + 二次计算 → UI/告警。


analyzer 模块群

oap-server/analyzer/ 下 7 个子模块:

图5

agent-analyzer · Trace 切分

把 trace segment 切分成 Source 对象(Service/Endpoint/ServiceInstance/ServiceRelation/EndpointRelation/DatabaseAccess…),交给 SourceReceiver 分发到 OAL 生成的 Dispatcher。这是 OAL 数据的主要来源。

核心是"监听器"机制:一个 trace segment 进来,按 span 在 trace 里的角色(First/Entry/Exit/Local/Segment)创建不同 listener,每个在遇到自己关心的 Point 时提取对应 Source(oap-server/analyzer/agent-analyzer/.../listener/)。AnalysisListenerAnalysisListener.java:24)是根接口,containsPoint(Point):34)声明关心哪些点。

  • FirstAnalysisListener:27):trace 第一个 span,提取 trace 级信息。
  • EntryAnalysisListener:27):入口 span(如 HTTP 入口),提取 endpoint 流量。
  • ExitAnalysisListener:27):出口 span(如调下游/DB),提取关系。
  • LocalAnalysisListener:27):本地 span,不产生跨服务关系。

具体实现:

  • SegmentAnalysisListener:49):把 segment 原始数据落盘(SegmentRecord)+ 采样决策。实现多个接口,一个监听器身兼多职。
  • RPCAnalysisListener:59):同时实现 Entry+Exit+Local,是提取服务间调用关系/拓扑的核心——出口侧建 ServiceRelation(client),入口侧建 ServiceRelation(server),这就是拓扑图的来源。
  • VirtualServiceProcessor/VirtualDatabaseProcessor/VirtualCacheProcessor/VirtualMQProcessor/VirtualGenAIProcessor.../listener/vservice/):处理"虚拟" span(对缓存/DB/MQ/AI 服务的访问),转成 CacheAccess/DatabaseAccess/MQAccess/GenAIProviderAccess 等 Source,供 core.oal 对应 OAL 规则聚合。为什么叫"虚拟"?因为这些 span 不是真实的服务调用,而是探针埋点记录的"对某类资源的访问",OAP 要把它映射成对应的 Source 类型才能进 OAL 聚合。

根因分析
“根因分析"指通过 trace 拓扑定位错误/慢的根因 span。RPCAnalysisListener 同时处理 Exit(client 侧)和 Entry(server 侧),能把一个跨服务调用在两侧分别建关系 metric,从而在拓扑上区分"谁发起、谁接收、错误归谁”。SegmentStatusAnalyzer.../listener/strategy/SegmentStatusAnalyzer.java)按策略判定 segment 是否错误/慢。采样由 TraceSegmentSampler/SamplingPolicy 控制——只持久化慢/错误 trace,全量 trace 只聚合成指标,这是存储成本和可观测性的权衡。

meter-analyzer / log-analyzer

见上文 MAL / LAL 节。meter-analyzer 是 MAL 的宿主,被 OTel/Envoy/Telegraf/Zabbix/log-mal-rules 等所有处理外部 meter 的 receiver 共用。log-analyzer 承载 LAL 编译器,同时是 log-mal-rules catalog 的宿主。

event-analyzer · 事件

EventAnalyzerModuleProvider 处理 SkyWalking “事件”(服务实例启动/停止、JVM GC 事件、告警事件),注册 EventRecordAnalyzerListener.Factory,事件经 listener 转成 Event record 存储。事件不是指标,走 RecordStreamProcessor,不走 OAL/MAL——因为事件是离散的、不需要时间窗口聚合的原始记录,强行走聚合反而丢失语义。

gen-ai-analyzer · AI/LLM 可观测

GenAIAnalyzerModuleProvideroap-server/analyzer/gen-ai-analyzer/.../GenAIAnalyzerModuleProvider.java:84OALEngineLoaderService.load(GenAIOALDefine.INSTANCE):85)加载 oal/virtual-gen-ai.oal,把 GenAI Provider/Model 的访问(GenAIProviderAccess/GenAIModelAccess Source)聚合成指标(token 数、成本、TTFT、延迟分位数)。它是一个专门��� OAL 触发器,配合 VirtualGenAIProcessor 工作——这是 OAL 可扩展性的体现:新可观测域只要定义 Source 类 + 写 .oal + 触发器,就能复用整套聚合/存储/查询链路。

ios-analyzer · iOS 端

监听器 IOSHTTPSpanListenerIOSMetricKitSpanListener.../analyzer/ios/listener/),处理 iOS SDK 上报的 HTTP span 和 MetricKit 数据,机制类似 agent-analyzer 但针对 iOS 客户端语义。

hierarchy · 层级聚合

用第五种小 DSL(层级规则表达式)把服务按名称/短名/命名空间匹配成"上下层级"(如 父服务.子服务),用于在 UI 里把分布式服务/进程聚合成层级树。

语法(oap-server/analyzer/hierarchy/src/main/antlr4/.../HierarchyRuleParser.g4)是 Groovy 闭包风格:{ (u, l) -> u.name == l.name }simpleExpression:46)。CompiledHierarchyRuleProvider.../config/v2/compiler/CompiledHierarchyRuleProvider.java:48)是 SPI 实现,被 HierarchyDefinitionService 通过 ServiceLoader 发现;HierarchyRuleClassGenerator.compile 用 ANTLR 解析 + Javassist 生成实现 BiFunction<Service, Service, Boolean> 的类,运行时由 HierarchyService 调这些匹配器建层级。

hierarchy 不属于四大 DSL
hierarchy 用 ANTLR + Javassist 编译,但它是独立的"匹配规则 DSL",只服务于"层级"这一个查询侧特性,不属于 OAL/MAL/LAL/MQE 四大 DSL。规则在 hierarchy-definition.yml(4 种内置模式:name/short-name/lower-short-name-remove-namespace/lower-short-name-with-fqdn)。


关键源码索引

OAL

  • 语法:oap-server/oal-grammar/src/main/antlr4/org/apache/skywalking/oal/rt/grammar/OALParser.g4:32root)、:56source)、:36disableStatement)、:214castStmt
  • 触发:oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/oal/rt/OALEngineLoaderService.java:58load)、:88(反射加载)
  • 引擎:oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/OALEngineV2.java:90start)、:134notifyAllListeners
  • 解析:oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/parser/OALScriptParserV2.java:99parse)、:123OALListenerV2
  • AST:oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/model/MetricDefinition.java:39
  • enrich:oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/generator/MetricDefinitionEnricher.java:96
  • 代码生成:oap-server/oal-rt/src/main/java/org/apache/skywalking/oal/v2/generator/OALClassGeneratorV2.java:143generateClassAtRuntime)、:196(Metrics 类)、:341(Builder 类)、:393(Dispatcher 类)
  • 接线:oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/StreamAnnotationListener.java:50
  • 落盘:oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/MetricsStreamProcessor.java:165(OAL)、:174(MAL)
  • 示例脚本:oap-server/server-starter/src/main/resources/oal/core.oaloal/disable.oaloal/virtual-gen-ai.oal

MAL

  • 语法:oap-server/analyzer/meter-analyzer/src/main/antlr4/org/apache/skywalking/mal/rt/grammar/MALParser.g4:47additiveExpression)、:64postfixExpression)、:79(扩展函数)、:126closureExpression
  • 加载:oap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/prometheus/rule/Rules.java:60loadRules)、:149(SnakeYAML)
  • OTel 触发:oap-server/server-receiver-plugin/otel-receiver-plugin/.../otlp/OpenTelemetryMetricRequestProcessor.java:218start)、:226loadRules)、:237new MetricConvert
  • 编排:oap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/MetricConvert.java:107(构造)、:132(Phase1)、:166(Phase2)、:306formatExp)、:325toMeter
  • 单规则:oap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/Analyzer.java:176prepare)、:218register)、:264analyse)、:325expression.run
  • 解析+生成:oap-server/analyzer/meter-analyzer/src/main/java/org/apache/skywalking/oap/meter/analyzer/v2/compiler/MALScriptParser.java:213parse)、:280MALExprVisitor);.../compiler/MALClassGenerator.java:150compile);.../compiler/MALMetadataExtractor.java:55
  • 存储注册:oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/meter/MeterSystem.java:113create
  • 示例脚本:oap-server/server-starter/src/main/resources/otel-rules/vm.yamlmeter-analyzer-config/java-agent.yaml

LAL

  • 语法:oap-server/analyzer/log-analyzer/src/main/antlr4/org/apache/skywalking/lal/rt/grammar/LALParser.g4:38root)、:43filterBlock)、:62parserBlock)、:111extractorBlock)、:186metricsBlock)、:227sinkBlock)、:235sinkStatement
  • 模块启动:oap-server/analyzer/log-analyzer/.../provider/LogAnalyzerModuleProvider.java:154start)、:160new MetricConvert
  • 解析:oap-server/analyzer/log-analyzer/.../compiler/LALScriptParser.java:86parse
  • 代码生成:oap-server/analyzer/log-analyzer/.../dsl/DSL.java:107compile)、:134evaluate)、:142execute
  • 运行时:oap-server/analyzer/log-analyzer/.../provider/log/listener/LogFilterListener.java:66(类)、:89build)、:114parse)、:539Factory.create
  • LAL→MAL 桥:oap-server/analyzer/log-analyzer/.../dsl/spec/extractor/MetricExtractor.java:67prepareMetrics)、:74submitMetrics)、:97toMeter
  • 示例脚本:oap-server/server-starter/src/main/resources/lal/nginx.yamllal/default.yaml

MQE

  • 语法:oap-server/mqe-grammar/src/main/antlr4/org/apache/skywalking/mqe/rt/grammar/MQEParser.g4:25expression)、:29addSubOp)、:35trendOP)、:39relablesOP)、:40aggregateLabelsOp)、:41sortValuesOP)、:56metric)、:66scalar)、:81topN
  • 入口:oap-server/server-query-plugin/query-graphql-plugin/.../mqe/rt/MQEExecutor.java:71execute)、:83(lexer)、:89parser.expression()
  • Visitor 基类:oap-server/mqe-rt/src/main/java/org/apache/skywalking/mqe/rt/MQEVisitorBase.java:64(类)、:110visitAddSubOp)、:178visitAggregationOp)、:404visitTrendOP
  • 具体 Visitor:oap-server/server-query-plugin/query-graphql-plugin/.../mqe/rt/MQEVisitor.java:65(类)、:110visitMetric)、:280readMetricsValues)、:311readLabeledMetricsValues
  • 结果模型:oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/mqe/ExpressionResult.java