MiroFish:开源数据流镜像与回放工具,让异常排查更快
2026/9/18 6:54:13 网站建设 项目流程

先交代个背景。我平时的工作离不开排查线上数据异常:每天早晨总有几条告警、订单状态不对、统计曲线突然塌一块。日志文件堆了几个G,但想在时间上把“问题前”和“问题后”的两段流逐条对上,靠人眼基本看不过来。后来我给自己写了MiroFish——一个开源的数据流镜像与回放工具。它不替代链路追踪,也不替代监控大盘,它只做一件事:把一段数据流水原样镜像下来,之后你可以像回放录像一样反复观察、逐条比对。

这套东西起初是给自己用的,后来陆续给同事用过几轮,大家反馈“比想象中有用”。尤其是想搞清楚“数据到底是在哪个环节被改掉的”“为什么测试环境复现不出线上效果”这类问题,MiroFish的镜像对比能力能省掉大半拉通排查的时间。如果你也在做后端服务、数据处理管道、支付对账、实时风控、消息队列相关的事情,或者只是经常被“同样的输入,结果却不一样”折磨,这篇内容基本就是为你准备的。

1. 项目定位与整体设计思路

1.1 为什么叫 MiroFish:数据流就像鱼群

名字是两个词拼出来的。Miro 在拉丁语系里有“镜子”的意思,Fish 就是鱼。我把数据流看成一条河里的鱼群,每条数据就是一条鱼,平时它在管道里游得很快,出了问题时你根本看不清是哪条鱼先“变异”的。MiroFish 做的事情,相当于在河道中间放了一面透明的玻璃墙:鱼游过去的时候,会被完整地照下来。照下来的影像可以反复看,也可以和另一段影像做逐帧对比。

这个比喻决定了整个项目的设计取向:我不指望它把所有问题都自动定位完,而是要它尽量完整、保真地还原“当时发生了什么”。所以它的核心不是复杂分析,而是镜像、存储、重放、对比这四个能力。你可以理解为:先给数据流拍一段“高清录像”,然后带着这段录像去复盘。

这也解释了为什么它的名字里带了 Fish 而不是 “MirrorFlow” 之类的词。鱼是有脾气的,河水也是混的,很多工具号称能帮你分析数据,但真正采集时不完整、回放时失真、对比时不对齐,最后还不如直接看原始日志。我不想再做这样一个“看起来很美”的工具。

1.2 整体架构:四面镜子各管一段

MiroFish 不是一个单体程序,我把它拆成了五个相对独立的模块,它们通过标准输入输出和 JSONL 文件互相连通。这样做的原因是,在实际使用中,数据源千差万别:有人从 Kafka 拉,有人读日志文件,有人直接抓 HTTP 请求,还有人只是想把本地一段 JSON 数组导进去做实验。如果把数据接入和核心分析强耦合在一起,每接入一种新数据源就要动核心代码,维护成本会迅速失控。

五个模块分别是:

  • capture:负责从各类数据源捕获原始事件,统一转成 JSONL 格式后写入存储。
  • store:负责落盘、索引和时间窗口管理,底层用 SQLite 加 WAL 模式。
  • mirror:负责对比两段镜像,支持按序号、时间戳或自定义主键对齐。
  • analog:负责重放一段镜像,可以按原始节奏,也可以按倍速或限速回放。
  • web:负责可视化,展示吞吐量、延迟分位数、异常规则命中情况。

这五个模块共享同一个配置文件和同一套事件格式。只要事件能变成 JSON 对象,MiroFish 就能处理。反过来,如果你接入了一个新数据源,意味着只需要给 capture 写一个适配器,其余部分不用动。这个分层思想是我从代理服务器和日志采集器的设计里学来的,代码不复杂,但边界清楚后,扩展和排查都快很多。

1.3 技术选型:为什么是 Python 加 JSONL 加 SQLite

先说实话:决定用 Python,不是因为 Python 性能最强,而是因为这类工具的瓶颈几乎永远不在语言本身,而在接入成本和分析灵活性。使用者的最常见操作是写一段清洗脚本、跑一个统计、看一眼分布,这种场景下 Python 的生态优势和开发速度太明显了。性能不够的地方,我用多进程和批量写入来补,实测单机处理每秒几万条事件没有压力,对绝大多数业务诊断和压测复盘场景已经足够。

事件格式我选了 JSON Lines,每行一个 JSON 对象。这个选择很朴素,但有几个实际好处:可以用 grep 直接搜关键字;可以用 jq 做临时分析;文件可以按时间滚动切割;即使 MiroFish 崩溃了,最后一条未写完的半行也不影响前面数据。相比二进制格式或者数据库直写,JSONL 牺牲了一点空间效率,但换来了巨大的可观测性和便捷性。

存储则用 SQLite,开启 WAL 模式。这里要解释一下:为什么不直接用 MySQL 或者 Elasticsearch?因为 MiroFish 的定位是“项目级的镜像工具”,不是“平台级的日志系统”。它存储的是短时间窗口内的结构化事件,通常就是几百 MB 到几个 GB 的量级。SQLite 单文件,备份方便,恢复也方便,镜像对比时可以直接把文件拷到另一台机器上分析。后来我在一个压测复盘场景里,直接把同事的 SQLite 镜像文件压缩后用即时通讯软件传过来,省去了搭服务的环节,这种轻量优势是重型存储给不了的。

2. 核心模块设计与配置解析

2.1 事件格式:所有模块的共同语言

模块之间不互相调用内部函数,而是约定了一种最小事件格式,MiroFish 内部管它叫 “MiroEvent”。一个事件最少得有四个字段:

{ "ts": 1710000000.123, "seq": 1024, "source": "order-api", "topic": "order.created", "payload": { "order_id": "A10086", "amount": 199.0, "channel": "wxpay" } }

字段含义如下:

字段类型说明
tsfloatUnix 时间戳,保留毫秒或微秒,尽量用事件发生时间而不是采集时间
seqint单调递增序号,作为同一数据源内的事件顺序依据
sourcestring数据来源标识,比如服务名、机器名、日志文件路径
topicstring业务类型,类似 Kafka 的 topic 概念,用于分组筛选
payloadobject业务原始数据,原样保留

为什么特意区分 ts 和 seq?因为网络传输和时间戳采集很可能导致两条事件到达顺序错乱,ts 可以告诉你“实际什么时候发生”,seq 可以告诉你“原始产生的顺序”。mirror 做对齐时优先用 seq,因为它在同一数据源内是可靠的;跨数据源对比时才用 ts 加容差时间窗。这点看着简单,但我早期踩过坑后才意识到,顺序语义是镜像对比的基石。

2.2 配置系统:一份 YAML 管全部

MiroFish 的配置文件是 YAML 格式,整个工具的所有行为都被收敛在这个文件里,避免每个人启动时打一长串参数。我最常用的一份配置大致长这样:

app: name: mirofish-demo storage_path: ./data/ timezone: Asia/Shanghai capture: input: kafka kafka: brokers: ["localhost:9092"] group_id: "mirofish-capture" topics: ["order.created", "order.paid"] offset: latest batch: size: 500 flush_interval_sec: 2 store: engine: sqlite sqlite: db_path: ./data/mirofish.db wal: true retention_days: 7 bucket: jsonl jsonl: dir: ./data/jsonl/ rotate_size_mb: 128 mirror: align_by: seq tolerance_ms: 200 analog: speed: 1.0 loop: false web: host: 127.0.0.1 port: 8745

配置里比较需要动脑的是 capture.batch 和 store.jsonl.rotate_size_mb。batch 控制捕获模块攒多少条再批量写入,size 太小时频繁写,数据库压力大;size 太大时一旦宕机,内存里未落盘的事件会丢得更多。我一般用 500 条一次,flush 间隔 2 秒兜底。rotate_size_mb 则是控制 JSONL 文件多大会滚动切割,128 MB 是我在 grep 和文件数量之间找到的平衡点,太大则打开慢,太小则文件碎片多。

2.3 mirror 镜像对比:最核心的逻辑

mirror 模块是整套工具里含金量最高的部分。它的用途是拿两段事件流做对比,找出“哪些事件在 A 里有而 B 里没有”“哪些事件两边都有但 payload 不同”。实现思路不复杂,但细节决定了结果到底可不可信。

对比过程分三步。第一步是定义基准侧和目标侧,通常基准侧是正常时段或线上黄金镜像,目标侧是问题时段或新版本测试结果。第二步是对齐,按 seq 相同则视为同一条事件来匹配,如果 seq 缺失或跨源对比,就按 ts 在 tolerance_ms 范围内找最近匹配。第三步是输出差异,分为 missing、extra 和 modified 三类。

这里的“modified”是最容易被忽视却也最有用的一类。它表示两个镜像里都存在同一条事件,但 payload 字段值发生了变化。比如订单金额从 199.0 变成了 19.9,mirror 会把变更字段单独列出来,并显示变更前后的值。我在一次联调里发现某个网关把金额单位从分转成了元,但只在特定渠道下触发,就是靠这个字段级差异定位的,单靠日志关键字根本搜不出来。

2.4 可观测性设计:不是只写文件就完事

如果数据只是落盘后等人来分析,那遇到线上问题时仍然太被动。所以 MiroFish 在 capture 和 store 之间加了一层实时统计:每秒记录事件数、平均大小、来源分布、延迟分位数等指标,并且这些指标本身也会作为元事件写入存储。这样当你事后打开 Web 面板时,第一眼看到的不是密密麻麻的原始数据,而是这段时间内的吞吐曲线和异常拐点,然后再“下钻”到具体事件。

这部分设计我参考了监控系统里“标签 + 指标”的思路。每个事件会自动附带几个标签(source、topic、node),统计时可以按任意标签组合做聚合。不需要预先声明所有维度,因为 JSONL 本身就是 schema-free 的。使用的时候,页面支持直接写一个简单的过滤表达式,比如 source == "order-api" and topic == "order.created",然后只看这部分数据的统计,这个体验已经被不少人夸过。

3. 从安装到跑通一份完整记录

3.1 环境准备与安装

MiroFish 目前以 Python 包的形式分发,支持 Python 3.9 及以上版本,依赖库尽量克制,核心只有 PyYAML、fastapi、uvicorn、pandas 这几个。安装命令很简单:

pip install mirofish

如果你所在环境不允许直接连外部包源,也可以从 Git 仓库克隆后离线安装:

git clone https://github.com/your-org/mirofish.git cd mirofish pip install -r requirements.txt python setup.py install

我建议用虚拟环境,尤其是公司服务器上同时跑着多个 Python 服务,依赖互相污染的情况遇过太多次。项目目录下创建一个虚拟环境,然后激活,再执行 pip install,这样最稳妥。安装完成后运行mirofish --version,能输出版本号就说明基础环境没问题。

3.2 用模拟数据打通全流程

不少初学者一上来就想接真实 Kafka 或真实数据库,结果配置复杂又看不到数据流转,费力不讨好。我的做法是先造数据,把整条链路跑通,再接真实数据源。

造数据我这里给你一个小脚本思路,大概作用是随机生成带有少量异常的订单事件流,存成六万行的 JSONL 文件,用来模拟一段二十分钟的线上流量:

import json import random import time base_ts = time.time() topics = ["order.created", "order.paid", "order.refund"] channels = ["wxpay", "alipay", "card"] with open("demo_events.jsonl", "w", encoding="utf-8") as f: for i in range(60000): ts = base_ts + i * 0.02 # 每秒50条 is_abnormal = (i > 30000 and i < 31500) # 模拟一段异常区间 payload = { "order_id": f"ORD{i:06d}", "amount": round(random.uniform(10, 500), 2), "channel": random.choice(channels), } if is_abnormal: payload["amount"] = round(payload["amount"] / 100, 2) # 模拟单位错乱 event = { "ts": round(ts, 3), "seq": i, "source": "demo-generator", "topic": topics[i % len(topics)], "payload": payload, } f.write(json.dumps(event, ensure_ascii=False) + "\n")

生成完成后,启动 capture 读取文件:

mirofish capture --config mirofish.conf.yaml --input file --file-path ./demo_events.jsonl

这里的--input file会覆盖配置文件里的 kafka 输入源,适合做本地演练。命令执行后应该看到每秒打印一条进度信息,显示已处理事件数和当前速率。跑完后,去data/jsonl/目录下看一眼,应该能看到按时间滚动的 JSONL 文件,同时 SQLite 数据库文件也已经生成。

3.3 启动 Web 面板查看统计结果

数据已经入库,接下来把 Web 面板跑起来,验证统计和图表的展示是否正常:

mirofish web --config mirofish.conf.yaml

浏览器访问http://127.0.0.1:8745,如果一切正常,第一屏会显示事件总吞吐曲线、各 topic 占比饼图、以及按 source 聚合的事件数排行。页面往下拉,能看到一个“关键分位数”区域,P50、P95、P99 延迟会以表格形式展示。我们这份模拟数据没有真实延迟字段,所以这里的“延迟”实际是事件到达间隔,也算是间接衡量速率稳定性的一种指标。

我特意在页面上做两个时间段对比的功能,比如选择“异常区间”和“正常区间”,页面会并排显示两个窗口的统计差异。这个功能不依赖 mirror 模块的精确对齐,主要给人快速浏览剩余异常位置的粗粒度差异,适合在复盘会议上先给团队一个直观印象,再进 mirror 做逐条精确对比。

3.4 实战演示:mirror 对比发现异常区间

现在进入重点环节,用 mirror 模块把正常区间和异常区间做一次完整镜像对比。操作流程是这样的:

先把 demo_events.jsonl 前面三万条作为基准镜像导出:

mirofish export --config mirofish.conf.yaml --where "seq >= 0 and seq < 30000" --output ./base.jsonl

再把三万一到三万六之间的数据作为目标镜像导出:

mirofish export --config mirofish.conf.yaml --where "seq >= 31000 and seq < 36000" --output ./target.jsonl

然后执行镜像对比:

mirofish mirror --base ./base.jsonl --target ./target.jsonl --align seq --output ./diff.json

命令结束后,打开 diff.json,会看到类似下面的输出:

{ "summary": { "base_total": 30000, "target_total": 5000, "matched": 0, "missing": 30000, "extra": 5000, "modified": 0 }, "details": [] }

为什么 matched 是 0?因为我把 seq 范围完全错开了,两条镜像根本没有交集,自然全算 missing 和 extra。这是演示时故意做的一个不太合理的对比,但它能说明一个问题:mirror 对比前最好确保两段镜像有重叠区间,否则结果会显示大量 missing。实际使用时,我会把基准镜像定义为“问题发生前五分钟”,目标镜像定义为“问题发生中五分钟”,两条镜像在时间上紧邻,seq 范围连续但不重叠,这样也能通过合理的时间窗口对比说明问题。

换一个更接近真实操作的例子,把目标镜像改为 seq 30000 到 33000,和基准镜像有一段重叠部分,执行对比后结果会输出三个区块。我在实际定位金额单位错乱时,就是发现payload.amount字段在目标侧数值整体缩小了 100 倍左右,mirror 在 details 区域明确标出了修改前后的值,几分钟内就确认是转换逻辑只在异常区间生效。

3.5 用 analog 模块回放复盘

最后演示重放能力。在实际故障复盘时,光看对比结果还不够,有时候需要在“故障现场”环境下再走一遍,比如验证修复后的程序是否能正确处理那段异常数据。analog 模块就是为了这件事准备的。

回放默认按原始时间节奏执行,也就是说,捕获时花了二十分钟的事件流,回放也要二十分钟。可以用--speed参数倍速播放,加速到 10 倍或 50 倍,也可以限速到 0.5 倍细细观察:

mirofish analog --config mirofish.conf.yaml --input ./target.jsonl --speed 10 --stdout

--stdout会把事件实时打印到控制台,配合 jq 可以边回放边过滤,比如只看 order.refund 类型的事件或者只观察存在异常特征的订单。回放还有一个隐含价值:它可以作为回归测试的数据源,让修复后的处理程序消费这段数据,看输出是否符合预期。这个用法一开始没在计划里,是团队一位测试同事提出的,后来成为我最常用的场景之一。

4. 常见问题与避坑经验

4.1 采集数据量比预期少,可能丢事件

很多人第一次接入自己业务数据后,发现 MiroFish 里的事件数明显少于日志里统计的条数,第一反应是工具写丢了。其实大多数时候不是写丢失,而是采集源消费语义的问题。

拿 Kafka 接入为例,消费者默认从 latest 开始消费,启动之前积压的消息根本不会进入 MiroFish;如果消费者组的问题导致分区分配不均,也会有个别分区始终没被轮到。还有一类情况容易忽略:业务程序在发送事件失败时可能自动重试,导致上游日志里有一条记录,但下游实际只成功收到一条,或者反过来收了两条。

排查这类问题,我总结了一个优先顺序:先看 capture 启动时打印的消费者组 offset 信息,确认起始位置;再看 Kafka 生产端的发送成功回调;最后把 MiroFish 记录的事件数和对账脚本统计的数量做分钟级对齐。大多数时候问题在起点,而不是在存储。

4.2 时间窗口对比结果总对不上

mirror 对比时另一大坑是时间窗口对不上。两个数据源分别部署在不同机器上,系统时间不一致,哪怕相差只有几百毫秒,在低延迟的业务场景下也会导致事件对齐不上,明明同一条事件却被判定为 modified。

解决办法有两个层面。配置层面,把 mirror 的tolerance_ms适当调大,比如从默认 200 调到 1000,给时钟偏差留缓冲。架构层面,更推荐在事件写入源头统一使用“上游服务处理时间”而不是本地采集时间,也就是说 ts 由生产者赋值,而不是由 MiroFish 的 capture 模块赋值。我后来在团队内部约定,所有业务事件里必须带上业务层时间戳,采集层不再覆盖这个字段,之后镜像对比的准确率提升非常明显。

如果两段镜像来自完全不同的系统,连 seq 都没有对应关系,那只能退而求其次,用 ts 加一些唯一键(比如订单号)做关联。这种跨系统对比的精度会低一些,但依然能定位出大部分明显的字段差异,我把它定位在“快速筛查”而不是“精确证明”的层级。

4.3 SQLite 文件越来越肥,请求变慢

SQLite 用久了之后文件膨胀是常见问题,尤其在高频写入场景下。MiroFish 默认只保留retention_days天内的数据,但删除数据并不会自动缩小文件大小,SQLite 只是把页面标记为空闲。当历史数据反复写入删除时,文件里会存在大量碎片页。

我一般建议做两件事:一个是定期执行 VACUUM 重建数据库,把空闲页回收掉;另一个是直接按天归档早期数据,把某一天之前的 SQLite 文件手动移走,需要分析旧数据时再挂载回来。这样生产库保持轻量,历史数据又不会丢。VACUUM 操作比较吃 I/O,建议在业务低峰期执行。

另外,开启 WAL 模式后,-wal文件会单独增长,正常情况下检查点会自动回收。如果观察到 wal 文件异常大,通常是长事务或者连接长期未关闭导致的,排查一下是否有客户端在事务里停留太久没提交。

4.4 性能优化心得:从“能用”到“好用”

MiroFish 在单机上的吞吐瓶颈通常发生在三个位置:Capture 的 JSON 序列化与反序列化、SQLite 的写入锁竞争、Web 查询时的大范围扫描。对前两个问题,我的经验是:事件尽量在进入 capture 前就是 JSON 字符串,capture 内部只解析必要的路由字段,payload 部分可以用原始字符串存储而不是重新序列化,这样能省掉很大一部分 CPU 开销。

SQLite 写入方面,开启 WAL 模式只是第一步,关键是使用事务批量提交。500 条一批提交一次,比每一条提交一次快一到两个数量级。实际压测中,我拿 100 万条模拟事件做过实验,在普通笔记本上,批量模式总耗时约 35 秒,平均每秒接近三万条,不再有“跑起来风扇狂转”的情况。

Web 查询慢则比较依赖索引设计。ts 和 seq 这两个字段一定要建索引,否则范围查询在几百万条数据上会明显卡顿。MiroFish 初始化时默认就会为事件表建 ts 和 seq 的联合索引,但如果你手工导入过数据,最好检查一下索引是否存在,不要假设它一定在。

4.5 几个容易忽略的“小技巧”

最后分享几个我在长期使用中总结的、文档里不会写的小技巧。第一个是关于镜像文件的命名规范。MiroFish 的 JSONL 滚动文件默认按时间命名,但我习惯在文件名里额外带上 source 标识,比如order-api-20250118-160000.jsonl。这样在磁盘上浏览时会非常直观,单独把某个服务的数据拿给同事分析时,也不会搞混来源。

第二个是模拟实时输入时的一个取巧方法。开发调试阶段如果没有真实消息队列,可以创建一个命名管道来模拟流式输入:

mkfifo /tmp/mirofish_pipe mirofish capture --config mirofish.conf.yaml --input file --file-path /tmp/mirofish_pipe

然后另一个终端往管道里写数据,capture 会像读实时流一样消费,效果接近真实环境。这个技巧我在没有测试环境的情况下用了很多次,简单有效。

第三个是给事件打“快照标记”。在关键时间节点或者发布操作时,我会给事件加一个快照标记,比如 payload 里写入一个snapshot: "v1.2.3-release"字段。这样以后做镜像对比时,可以一次性筛出某个发布版本的所有事件,直接当成一个黄金镜像来使用。这个习惯帮我节省了大量准备对比基线的时间,也让每次发布的变更影响范围变得可视。

我个人在实际操作中的体会是,MiroFish 这类工具的价值,往往不是它有多“智能”,而是它让你拥有了“回看数据现场”的能力。线上问题最怕的就是错过现场,有了完整镜像之后,很多原本需要靠猜的问题都可以变成基于证据的定位。如果你经常被数据不一致问题困扰,不妨从今天开始,给你的一条核心数据流做一份镜像存档。等到某天凌晨被告警叫醒时,你会感激自己这个决定的。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询