流式事件不是只能进规则引擎:用 BlazeRules 把日志、Kafka 与决策证据放到同一条数据路径
线上系统要对事件做判断,常见做法是把规则散落在应用代码、SQL、流处理任务和告警平台里。订单风控、支付异常、日志分流、数据质量检查看似不同,最后都会遇到同一个问题:一条事件为什么被拦下、由哪版规则决定、坏数据去了哪里?
BlazeRules提供的是另一种边界明确的实现:它以 YAML 规则描述判断条件,把输入收集成微批(microbatch)后执行,并把紧凑的决策记录或死信记录输出。项目是 Apache-2.0 许可的 C++ 工程,同时提供 Python 包、原生 CLI、本地接入 Agent 与只读 Dashboard。它不替代业务系统的事务逻辑;更适合放在事件进入后续处理之前,承担“可声明、可复放、可观察”的判断层。
先把“规则”与“传输”拆开
规则系统容易失控,往往不是条件表达式太复杂,而是把输入协议、领域判断和输出动作绑死了。BlazeRules 的结构可拆成四段:接入适配器收集事件;解码器把 NDJSON、Arrow、Avro、Protobuf 等输入变成批;规则引擎按 YAML 判断;输出层写入决策、摘要或死信数据。
这意味着同一份规则不必只服务 HTTP 接口。官方文档列出的接入面包括文件、stdin、HTTP、Kafka、Debezium CDC、Arrow IPC、Parquet、CSV、Avro 与 Protobuf。实际部署时,先问“事件从哪里来、结果由谁消费”,再决定接入方式;不要因为规则文件是 YAML 就把它误当作配置管理工具。
例如,离线核对一批 NDJSON 事件可以直接使用 CLI:
blazerules eval \ --rules rules.yaml \ --input ndjson \ --path events.ndjson \ --output grouped-decisions
这里的 eval 很适合做回放:在上线新规则前,用历史事件比较命中情况;变更后保留规则版本和输出摘要。它比只看“规则通过了语法校验”更有价值,因为规则的真正风险通常来自字段缺失、类型变化与条件组合,而不是 YAML 是否能解析。
用一份小规则建立可复核的决策
以支付事件为例,规则不需要一开始就做成复杂的风险模型。先把一个能解释的条件固定下来:金额高于阈值且设备类型为模拟器时进入人工复核。
- id: high_risk_payment
action: review
conditions:
and:
- field: amount
op: gt
value: 1000
- field: device_type
op: eq
value: emulator
YAML 只是规则载体;字段名、比较运算符和实际事件 schema 才是契约。投入生产前,应将 amount 的单位、空值语义、device_type 的枚举来源写清楚。否则“1000”可能在不同生产者中代表元、分或浮点金额,规则本身再快也只会稳定地输出错误结论。
项目支持在首次批处理时根据规则引用字段推断 schema,也允许显式 schema 用于更严格的类型控制。对金额、时间戳、国家码这类关键字段,建议将 schema 视作接口的一部分:版本化、在 CI 中校验,并用脱敏样本回放。不要把自动推断当成数据治理的替代品。
在线接入:吞吐、确认与死信必须分别设计
当应用需要实时发送日志或事件时,可以运行本地 Agent,暴露 HTTP 接入点:
blazerules_agent \ --rules rules.yaml \ --input http \ --host 127.0.0.1 \ --port 9480 \ --batch-size 4096 \ --ack-mode durable \ --output ndjson \ --output-path decisions.ndjson
随后,生产者把 NDJSON 发送到 /v1/logs。durable 确认模式的含义是输出写入后再响应;项目也提供在规则求值完成后响应的 evaluated 模式。两者不是谁更快谁更好:前者缩小“调用方收到成功但结果尚未落盘”的窗口,后者适合接受异步 sink 的低延迟链路。应按补偿能力、消息重试方式和结果存储可靠性选择。
坏记录也不应悄悄消失。BlazeRules 的文档说明,死信记录会保留错误码、出错列名和解析消息。把死信单独写成 NDJSON,再将其纳入告警与回放流程,能把“规则没有命中”与“事件根本无法解析”区分开。前者是业务判断,后者是生产者或 schema 契约故障;混在一起会误导排障。
把规则当作发布物,而不是线上临时开关
要让这条路径能长期维护,规则文件还需要像代码一样进入交付流程。建议每次变更都至少包含三类测试样本:应命中的正例、不应命中的反例、以及字段缺失或类型不符的边界例。对每个样本,记录预期的 action、命中规则 ID 和是否进入死信;部署前用同一份 CLI 在 CI 中回放。这样,当某个条件从 gt 改为 gte,团队看到的不是抽象的 YAML diff,而是哪些历史事件会改变去向。
还应明确规则版本的生命周期。一个可操作的做法是把规则文件随构建产物发布,并让输出中的 ruleset_version 对应 Git tag 或不可变构建号。消费决策日志的下游表应保留事件 ID、规则版本、命中 ID、决策时间和原始事件位置(例如 Kafka topic/offset 或对象路径)。发生争议时,复核者才能重建“当时输入的是什么、当时运行的是哪份规则”,而不是用今天的规则重算昨天的事件。
这也解释了为什么不要在 Dashboard 上直接把统计数字当作治理结论。某条规则命中率突然下降,可能是攻击减少,也可能是生产者换了字段名、某个解码器开始把记录送进死信,或者新版本规则改变了优先级。监控应把总输入量、可解析量、死信量、各规则命中量和输出写入延迟放在一起看;只看命中数会把数据管道故障误读成业务趋势。
批处理快,不等于可以跳过容量验证
另一个常被忽略的边界是规则优先级与可解释性。两条规则可能同时命中同一事件:一条要求人工复核,另一条要求直接阻断。无论引擎的默认组合策略是什么,团队都应在规则评审时把冲突处理写成可测试的预期:是按文件顺序、显式优先级,还是输出全部命中并交给下游仲裁。对涉及资金、权限或合规的流,最好在回放报告中同时列出“最终决策”和“所有匹配规则”,避免后来只看到一个动作却无法解释其竞争条件。
项目的卖点之一是 C++/SIMD 执行路径,并提供运行时分派的 AVX2、AVX-512 或 NEON 后端;但这不构成某个业务场景的性能承诺。真实吞吐还取决于输入解码、规则复杂度、模型调用、磁盘或对象存储输出,以及队列背压。
因此,容量测试至少要覆盖三组数据:正常事件、字段缺失或类型错误事件、规则高命中事件。还要分别观察 HTTP 接收、求值与 sink 队列。官方说明这些队列可独立设上限;正确做法不是先增大所有深度,而是先找出阻塞发生在接入、计算还是写出阶段。
如果规则需要模型评分,BlazeRules 可在条件中引用已注册的 ONNX 模型。但应保留纯规则的基线链路:模型不可用、特征漂移或延迟异常时,到底是降级为人工复核、拒绝事件,还是停止消费,必须由业务策略明确规定,不能靠默认行为猜测。
适合哪些团队
BlazeRules 适合已经有结构化事件、又不想为每类数据流重写判断逻辑的团队:例如把 Kafka 中的交易事件做预筛,把 CDC 变更做合规校验,或把容器日志转成结构化决策记录。它尤其适合需要回放和审计的场景,因为规则文件、输入样本、规则版本与决策输出可以共同构成复核材料。
相反,若数据本身仍是无定义的自然语言日志,或团队尚未确定字段契约,先做事件结构化会比引入规则引擎更重要。规则引擎能统一执行,不能替你决定字段含义;Dashboard 能展示结果,也不能证明阈值合理。
一个务实的起点是:选一条低风险事件流,用一条可解释规则跑离线回放;确认字段、误报与死信处理后,再接入 HTTP 或 Kafka。这样引入的不是“又一个规则平台”,而是一条能把事件判断、错误隔离和证据留存连起来的数据路径。