ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

ruflo:用 Rust 打造轻量流式日志处理管道

ruflo:用 Rust 打造轻量流式日志处理管道 1. 为什么会有 ruflo8.7GB 日志把我逼到写了个工具先说结论上周我又被一份 8.7GB 的网关日志虐了一遍和过去几年一样我第一反应是 grep、awk、jq 轮流上。等我把 awk 脚本从 2 分钟优化到 40 秒忽然意识到这有点不对劲——我只是想按分钟查一下 5xx 请求分布为什么要在一台 16G 内存的电脑上跟一堆管道和临时文件较劲于是我把踌躇很久的 ruflo 真正写完并开源了。ruflo 是一个用 Rust 写的流式处理管道工具名字取自 Rust Flow定位很朴素让日志、文本流、JSON 行这类数据能在一个轻量进程里完成过滤、转换、窗口聚合和转发。它既可以当命令行工具用也可以作为 Rust 库嵌入到项目里。你可以把它想成一个能自己拼管道的 jq或者一个部署成本几乎为零的轻量 Fluent Bit。这篇文章不打算写成一份官方文档我会把 ruflo 的设计取舍、核心抽象、完整跑通示例以及开发过程中踩过的背压、线程调度、JSON 解析三个大坑全部交代一遍。如果你也在做日志处理、数据管道或者单纯好奇 Rust 流式处理能压到什么程度这篇应该能给你一些不太一样的参考。1.1 一个具体的问题现场那天晚上的原始需求就一句话从 8.7GB 的 nginx 访问日志里把 5xx 请求按分钟统计出来看看到底是哪个时段出了问题。日志长这样服务器是 nginx 默认 combined 格式加了一些字段10.12.3.4 - - [12/Jan/2025:14:32:11 0800] POST /api/order/create HTTP/1.1 502 183 - Go-http-client/2.0 0.023 10.12.3.9 - - [12/Jan/2025:14:32:12 0800] GET /api/user/info HTTP/1.1 200 2453 - curl/8.4.0 0.011 10.12.3.7 - - [12/Jan/2025:14:32:13 0800] POST /api/payment/notify HTTP/1.1 503 182 - Java/1.8 0.108文件是 6200 万行单行平均 140 字节。我当时用 awk 做了一个过滤加按分钟聚合脚本跑完花了 1 分 52 秒峰值内存 380MB因为 awk 把每一分钟的状态码计数都放进了关联数组最后再统一遍历输出。这不是我第一次被这种操作烦到。之前处理 Kafka 消费延迟问题时我对着几 GB 的 consumer 日志做同样的事先 grep 一把再 awk 二次处理中途因为管道缓冲把内存打满了。工具不是没有而是没有一个刚好卡在本地排查 流式处理 低内存这个位置上的东西。1.2 为什么现有组合拳总差一口气我把常用方案都列了一下各有各的别扭grep sort awk快但是多步组合要反复串进程管道间缓冲不可控数据量大一点就很容易在某个中间环节把内存打满。而且每步都在重复读取同一个文件实际效率并没有想象中高。jq 管道JSON 处理能力很强但它是面向单条数据的处理海量 JSON Lines 时要先转换成 JSON 数组再处理吞吐和内存都不太够用跑 6200 万行很容易吃到 1GB 以上。Fluent Bit / Vector功能完整但配置体系重二进制体积大更多是为服务器采集场景设计的。我只是想在终端里快速看一批文件把它们拉起来有点杀鸡用牛刀。Logstash更不用说JVM 一启动就占几百 MB 内存本地排查场景直接劝退。我想要的工具必须满足四个条件单文件二进制、内存稳定在几十 MB、支持流式处理不落盘、既能写复杂配置又能一句话跑起来。市面上的工具没有一个同时满足这是 ruflo 出现的最直接原因。1.3 ruflo 的定位和名字由来名字取 Rust Flow但我更喜欢的解释是ruff 加 flow暗示它应该像 ruff 这个 Python linter 一样快、一样清爽。ruflo 的核心理念是一切数据都是流一切处理都是管道。具体来说它做三件事定义了一个统一的流式管道抽象Source 负责读、Operator 负责转换和聚合、Sink 负责写出。提供了一套类似 jq 的表达式语法适合终端快速操作也提供配置文件方式适合多阶段复杂管道。暴露了 Rust 库 API让开发者能在自己的项目里直接嵌入流式处理能力不需要额外起一个服务。单文件 release 构建目前是 2.9MB没有运行时依赖放到任何服务器上都能跑。这个体积带来的好处是排查问题时可以直接拷到生产机器上不污染环境用完就删。2. 核心抽象Source、Operator、Sink 三层模型ruflo 的整体设计其实只有三层但把这三层想清楚花了我不短的时间。早期版本我用过一个流水线 插件的模型每一阶段都是一个异步 trait 对象灵活是灵活但写起来非常啰嗦用户要理解什么是阶段阶段之间怎么传数据这些概念。后来我把模型简化成 Source、Operator、Sink 三段每一个输入进来的数据记录都在这三段之间流转。这个设计参考了 Unix 管道的哲学每个程序只做一件事通过标准 IO 连接起来。ruflo 只是把这件事搬到了内存里并且让每个环节都能独立控制并发和缓冲。2.1 管道模型和消息结构管道里的数据记录不是裸字节而是统一的Record结构pub struct Record { pub key: OptionString, pub timestamp: DateTimeUtc, pub fields: HashMapString, Value, pub raw: Bytes, pub source_offset: u64, }每条记录都保留原始字节和源偏移量这为后面排查问题提供了很大的便利——无论管道中间转了多少手最后都能追溯到它来自源文件的哪一行。timestamp在 Source 阶段就被解析出来后续所有窗口聚合都依赖这个字段而不是等 Operator 再去解析否则每个阶段重复解析时间戳的成本会非常可观。管道内部的连接默认使用有界队列队列长度 1024这是我的默认值可以通过参数调整。有界是必须的原因后面讲背压的时候会详细展开。2.2 Source 与 Sink 的对称设计Source 和 Sink 的接口是对称的这带来的最大好处是任何 Source 的输出都能接到任何 Sink 上组合方式是笛卡尔积的而不是为每种场景写死一个入口。目前内建的 Source 有stdin从标准输入逐行读取适合跟tail -f或journalctl配合。file顺序读取本地文件读完即止。tail类似tail -f持续监听文件末尾的新增内容。kafka从 Kafka topic 消费属于扩展能力默认不开启。Sink 有console输出到终端支持纯文本和 JSON。file写入本地文件。http按批 POST 到指定地址。null性能测试时用来丢弃数据。对称接口的具体落地是 Rust 的 trait#[async_trait] pub trait Source: Send Sync { async fn run(mut self, ctx: mut SourceContext) - Result(); } #[async_trait] pub trait Sink: Send Sync { async fn write(mut self, batch: [Record]) - Result(); }2.3 Operator 的三种类型Operator 是管道里最核心的部分我把它分成三种分别解决不同的问题类型作用典型场景Filter决定记录是否继续往下走status 500、level ERRORTransform修改记录内容映射到新字段正则提取、JSON 解析、字段重命名Aggregate按窗口和 key 做聚合每分钟计数、求均值、去重Filter 最轻不需要维护状态进来的记录要么放行要么丢弃。Transform 稍微重一点但本质上是逐条映射。Aggregate 是最复杂的因为它需要维护窗口状态——流式数据没有结尾的概念你永远不知道某个窗口的数据是不是已经全部到齐了。ruflo 的 Aggregation 采用 watermark 机制Source 会周期性下发 watermark 时间戳标识截至这个时刻之前的数据我都已经发完了。窗口聚合器收到 watermark 后才能安全地关闭并输出那些时间已经完整覆盖的窗口。这个机制是从 Flink 学来的只不过在轻量场景下实现起来大幅简化不需要分布式协调。2.4 管道表达式和配置文件的取舍ruflo 提供两套操作方式。终端场景用表达式把一段管道用|串起来类似 jq 的链路ruflo from(stdin://) | parse(nginx) | where status 500 | count window1m | to(console://)配置文件适合多阶段、参数复杂的管道。同样的逻辑用 TOML 表达[pipeline] name error-rate [source] type stdin [[stages]] type parse format nginx [[stages]] type filter expr status 500 [[stages]] type count window 1m keys [status] [sink] type console output json两套表达在内部会编译成同一个管道对象。表达式语法是给快速排查用的配置文件是给可复用任务用的互不冲突也没必要合二为一。3. 从零跑通 ruflo环境准备与两个最小示例这一节我按自己写文档的习惯来只讲能跑的、不绕弯子的路径。如果你已经装了 Rust 工具链整个过程应该不超过十分钟。3.1 安装方式ruflo 发布在 crates.io 上直接安装即可cargo install ruflo安装完成后验证一下ruflo --version如果想要最新的 main 分支也可以从源码编译git clone https://github.com/ruflo/ruflo.git cd ruflo cargo build --releaserelease 构建出来的二进制在target/release/ruflo拷到任意 Linux 或 macOS 机器上都能直接跑不需要额外的动态库。3.2 嵌入式 API 示例五分钟接入 Rust 项目下面这个例子演示在 Rust 项目里把 ruflo 作为库使用。场景是读取一个 nginx 访问日志文件过滤出 5xx 请求按分钟统计最后输出 JSON 到终端。use ruflo::prelude::*; #[tokio::main] async fn main() - Result(), Boxdyn std::error::Error { let pipeline Pipeline::builder() .source(FileSource::open(access.log)?) .stage(Parse::nginx()) .stage(Filter::expr(status 500)?) .stage(Count::window(1m).by(status)) .sink(ConsoleSink::json()) .build(); pipeline.run().await?; Ok(()) }这段代码里值得注意的一个设计是Filter::expr接受的是一个字符串表达式而不是编译期闭包。这意味着过滤器可以完全由配置文件或运行时参数驱动不需要改代码重新编译。代价是每次匹配都会解析一遍表达式但实际测试下来在流式场景下这个开销可以忽略因为表达式本身很简单不是热点路径。如果你想要极致的性能也可以直接传闭包绕开表达式解析.stage(Filter::new(|record| { record.fields.get(status) .and_then(|v| v.as_u64()) .map_or(false, |s| s 500) }))闭包方案适合画性能上限表达式方案适合写灵活业务两个 API 都保留。3.3 CLI 表达式示例像 jq 一样搞定日志流命令行场景我用一个更贴近日常操作的例子。假设你在排查一个线上问题需要实时观察 Nginx 日志里非 200 状态码的请求并按状态码每分钟计数tail -f /var/log/nginx/access.log | ruflo from(stdin://) | parse(nginx) | where status ! 200 | count bystatus window1m | to(console://)输出会长这样{window_end:2025-01-12T14:33:00Z,status:502,count:12} {window_end:2025-01-12T14:33:00Z,status:503,count:8} {window_end:2025-01-12T14:33:00Z,status:499,count:3}这里parse(nginx)会把一行日志解析成带命名字段的 Record之后where status ! 200才能引用status字段。窗口聚合输出的是窗口结束时间而不是开始时间这是为了避免还没到时间就有部分结果的误导。3.4 用配置文件跑一次完整的多阶段管道复杂任务建议用配置文件。我本地排查用过一个更完整的管道读取 8.7GB 的日志文件过滤 5xx提取用户 ID按用户 ID 聚合请求数最后把 Top 20 结果写到一个 CSV 文件里。配置文件如下[pipeline] name top-error-users [source] type file path /data/access.log [[stages]] type parse format nginx [[stages]] type filter expr status 500 [[stages]] type transform script record.fields[user_id] extract(record.raw, user_id(?Puid[0-9]), uid); return record; [[stages]] type count window 1h keys [user_id] [[stages]] type top key count limit 20 [sink] type file path /tmp/top_error_users.csv format csv运行ruflo run --config pipeline.toml这个例子覆盖了五种 stage 的组合。有一点要说明transform里的script目前是一个受限的表达式子集不是完整 JavaScript只支持字段读取、正则提取、简单的字符串和数值操作。设计成这样是为了保证执行链路不引入 JIT 或者脚本引擎维持内存和 CPU 的可预测性。4. 实测表现跟 grep awk jq 正面比一次工具好不好跑一次数据就知道了。这一节我把 ruflo 和传统组合拳放在同一个数据集上做了对比测试过程和结果都记录在这里。4.1 测试环境与数据集测试机器是 MacBook ProApple M1 Pro16GB 内存macOS 14。数据集是用下面命令生成的 6200 万行 nginx 风格日志ruflo gen --lines 62000000 --pattern nginx /tmp/access.log生成完看一眼体积$ du -sh /tmp/access.log 8.7G /tmp/access.log任务统一设定为过滤出状态码大于等于 500 的行按分钟聚合输出结果 JSON。这个任务覆盖了过滤和聚合两个核心操作比较有代表性。4.2 耗时与内存对比方案耗时峰值 RSS备注grep awk 组合2m08s380MBawk 数组累积 最终遍历输出jq 管道先转 JSON Lines6m1.1GB每行解析成Value树ruflo 单线程38s22MB数字恒定为最佳测试结果ruflo 4 线程11s24MB单机本地文件场景推荐配置grep awk 的 2 分 08 秒里grep 本身只花了 8 秒剩下的时间几乎全在 awk 的关联数组和最终遍历上。jq 方案因为要先把 nginx 格式解析成 JSON再交给 jq 处理中途还生成了一个 12GB 的中间文件所以拖到了 6 分钟以上。ruflo 单线程 38 秒的成绩已经让我满意4 线程 11 秒的结果则是意外之喜。内存从头到尾稳定在 24MB 左右没有明显的增长曲线这正是有界队列加批处理的效果。4.3 干活时不只靠快输出稳定性和可观测性跑得快只是其一我更看重的是在大文件场景下输出的稳定性。awk 脚本在内存到达 380MB 时CPU 时间大量花在哈希表扩容和重哈希上jq 更明显内存接近 1.1GB 时系统开始频繁换页整个命令的响应速度变得非常不稳定。ruflo 在运行中可以打印内部指标这是排查问题时的救命功能ruflo run --config pipeline.toml --stats输出类似[14:32:10] source: 1245000 lines, queue_depth: 32, mem: 20MB [14:32:20] source: 2482000 lines, queue_depth: 45, mem: 21MB [14:32:30] source: 3718000 lines, queue_depth: 40, mem: 22MB队列深度和内存这两个指标能非常直观地反映管道是否健康。如果队列深度持续上涨说明下游消费速度跟不上这时候加线程或优化 Sink 才是正确的方向如果队列深度很低但 CPU 打满则瓶颈在上游解析逻辑加线程没有意义。5. 踩坑实录背压、正则线程池、JSON 解析这三大硬骨头这部分是全文最值得看的内容。ruflo 开发过程中我遇到的几个问题每一个都不是文档里能直接查到的答案而是需要从现象倒推根因的排查过程。我把它们完整记录下来希望能帮你避开同样的坑。5.1 背压失效内存从 30MB 飙到 800MB 的排查链路第一个版本里我在 Source 和 Operator 之间用的是tokio::sync::mpsc::unbounded_channel。理由是简化逻辑反正总有办法控制流量。这个决定在第一版测试用例里完全没问题因为测试用的数据量只有几百条日志。直到我把 ruflo 接到一个真实的 HTTP Sink 上。这个 Sink 通过公网 POST 数据到采集服务服务端吞吐有限每秒只能接收 2000 条左右。而 Source 从本地文件读取的速度是每秒几十万条两边的速度差了两个数量级。当时我看到的现象是程序刚启动 30 秒内存从 30MB 平稳爬升到了 800MB并且还在继续往上涨CPU 却没有跑满。我的第一反应是怀疑 JSON 解析或者正则匹配有什么内存泄漏但用ruflo stats一看队列深度一栏触目惊心queue_depth: 4310512四千多万条记录堆积在 channel 里每条都带着完整的 HashMap 字段和原始 Bytes内存不被吃掉才怪。这让我意识到问题不在解析逻辑而在管道本身。排查到这一步根因已经清楚了unbounded_channel 没有任何背压机制生产者只管往队列里塞消费者处理不过来时内存就会被无限消耗。修复方案很直接把 channel 换成有界队列let (tx, rx) tokio::sync::mpsc::channel::Record(1024);然后 Source 端使用reserve模式发送队列满的时候自动等待let permit tx.reserve().await.map_err(|_| SourceError::Closed)?; permit.send(record);这样当 Sink 处理慢时背压会沿着管道一路传导回 SourceSource 会暂停读取文件文件描述符停留在当前偏移量等下游恢复后再继续。内存也稳定了下来。提示写流式处理程序时unbounded_channel是一个极具迷惑性的陷阱。测试数据小的时候永远发现不了问题一旦真实数据上来内存就是第一个暴雷的地方。默认情况下一律用有界队列只有在明确知道生产者速度一定低于消费者时才考虑无界。5.2 spawn_blocking 做正则匹配CPU 只吃掉一半的诡异现象第二个坑出现在我优化正则过滤性能的时候。ruflo 最初把 Transform 阶段放在了 Tokio 异步上下文里运行但正则匹配是 CPU 密集型任务放到 async 上下文里会阻塞 Tokio 的工作线程让整个管道的并发能力下降。咨询了一位做 Rust 服务端的朋友后我决定用tokio::task::spawn_blocking把正则匹配丢到阻塞线程池里跑。跑完 benchmark 我愣住了吞吐不但没有上升反而从 220MB/s 降到了 90MB/s而且top显示 CPU 使用率只有 30% 左右怎么看都不正常。排查时我分了三步用top和perf top观察进程状态发现大量时间花在epoll_wait和线程上下文切换上。把is_match替换成一个恒定返回true的闭包吞吐依旧是 90MB/s说明问题不是正则本身而是spawn_blocking这个调度方式。统计每次spawn_blocking调用的开销发现对于一个只需要几微秒就能完成的正则匹配来说调度开销反而占了主导。原因很清晰spawn_blocking是为那些可能耗时很长的阻塞任务设计的它会把任务从异步运行时转发到独立的阻塞线程池中间涉及跨线程通知、任务队列和唤醒。这些机制对于秒级的任务来说是合理的但对于微秒级的正则匹配每次匹配都要付出比匹配本身还贵的调度成本反而拖垮了吞吐。最终方案是把 Transform 阶段彻底从 Tokio 里剥离出来改用独立线程池配合 work-stealing 调度。我把这一层的实现换成了 rayonlet pool rayon::ThreadPoolBuilder::new() .num_threads(num_cpus::get()) .build()?; pool.install(|| { records.par_iter() .map(|record| transform(record)) .collect() });为了兼顾顺序ruflo 会先对记录按键做哈希分片相同 key 的记录一定进入同一个 rayon 线程内处理。这样既拿到了多核并行又不会破坏单 key 内部的时间顺序。调整之后同样的数据集上正则过滤吞吐从 90MB/s 恢复到了 280MB/s。这个坑让我学到的教训是在 Rust 异步生态里不是所有 CPU 密集任务都适合丢给spawn_blocking要区分偶尔阻塞和高频计算两种场景前者用阻塞线程池后者用专用计算线程池。5.3 serde_json::Value 拖垮吞吐换成流式反序列化第三个坑是 JSON 性能问题。ruflo 支持对 JSON Lines 格式的数据做解析第一版实现很简单直接调serde_json::from_slice::Value把整条 JSON 解析成一个递归的 Value 树。代码很干净但性能出来之后我有点不敢相信。测试数据是 6200 万行 JSON每行 150 字节左右解析吞吐只有 80MB/s比 nginx 格式的解析慢了三倍。更离谱的是用 heaptrack 跟踪内存分配时发现程序执行期间进行了至少几亿次小内存分配和释放。问题出在Value这个枚举上。serde_json::Value是递归类型解析时每遇到一个对象或数组都会在堆上分配新的节点遇到字符串又要分配一次 Vec。一行 JSON 里如果嵌套三层可能触发十几次甚至几十次堆分配。6200 万行乘以几十次分配次数就上亿了。修复思路很明确不要解析出完整的Value树直接让 serde 把需要的字段反序列化进一个紧凑的结构体#[derive(Deserialize)] struct RawRecord { status: u16, user_id: Optionu64, request_time: f32, } fn parse_json_line(line: [u8]) - ResultRawRecord { let mut de serde_json::Deserializer::from_slice(line); Ok(RawRecord::deserialize(mut de)?) }这样一个结构体里全都是栈上或紧挨着的字段不再有递归堆分配吞吐直接提升到了 250MB/s内存占用也降了一个数量级。如果你也要处理大量 JSON 数据我建议记住这条原则能用结构体就别用Value能用流式Deserializer就别用from_slice。Value是给那些字段结构完全不确定的场景准备的一旦你明确知道要拿哪些字段就应该给它一个确定的类型。5.4 CtrlC 丢最后一批数据优雅退出的实现细节最后一个坑也是 ruflo 从玩具走向可用工具的分水岭终端里按 CtrlC最后一批仍在聚合窗口里的数据会直接丢。现象非常隐蔽。我在测试 HTTP Sink 时管道每 1 分钟推出一批聚合结果我掐着表在窗口关闭前 5 秒按下 CtrlC发现接收端缺了最后一分钟的数据。如果管道是一条长跑任务这不算大事但对于处理本地大文件的任务最后一批结果往往是最重要的——因为它覆盖的恰好是文件末尾最新的数据。排查时我意识到默认的 CtrlC 行为对异步程序来说非常粗暴SIGINT信号直接把整个进程终止了Tokio 的运行时连清理缓冲区的机会都没有。任何还停留在 channel 里或 sink 缓冲区的数据都随着进程退出消失了。修复方案是为管道注册一个优雅退出流程用tokio::signal::ctrl_c()监听 SIGINT。收到信号后先将 Source 切换为 draining 状态不再读取新数据。让管道中已有的记录继续往下游流动每个 Operator 和 Sink 最多等待 5 秒完成 flush。flush 完所有缓冲后进程才真正退出。核心代码大概是这样的let shutdown async { tokio::signal::ctrl_c().await.expect(failed to listen); tracing::info!(received ctrl_c, draining pipeline...); ctx.shutdown().await; }; tokio::select! { result pipeline.run() result?, _ shutdown {} }注意这里的select!很关键——不能用pipeline.run().await之后再监听信号那样信号排队的时候管道还在跑永远不会执行到监听代码。这个改动上线之后我又专门写了一个测试脚本模拟窗口即将关闭时触发退出的场景重复跑了二十次确认每次都能完整输出最后一批窗口结果才把这个问题关闭。如果你在自己的异步程序里也碰到类似现象我的建议是先别急着加各种复杂的优雅停机框架先想清楚退出时有哪些数据还散落在各处按 Source 停读、管道排空、Sink 落盘的顺序把清理流程写出来往往就够了。整个 ruflo 开发下来我最大的体会是流式处理工具的难点从来不在怎么写解析和过滤而在怎么管住内存、怎么控制并发、怎么让进程在任何时刻退出都能保证数据不丢。这三个问题你在处理 10MB 日志时永远遇不到一旦数据量到了几十 GB它们会一个接一个冒出来而且每个都要靠实打实的测量和排查才能定位。现在每次给 ruflo 加新功能我都会顺手跑一遍背压场景的压测和退出场景的回归用例这两个测试已经成了项目里最宝贵的资产。
返回列表