OASIS BLOG
← 回到七日冲刺← D6 笔记 · D8 笔记 →
写给没碰过股票、期货、合约的 Java 程序员

D7 零基础速成笔记

最后一天不学新产品,而是把前六天做的东西主动弄坏:重复消息、超时、断线、热点、崩溃,看它在哪里出事、怎么防住。页面里带下划虚线的词,点一下就能看解释。

怎么读按顺序读,一共 12 节,大约 2–3 小时。第 3–7 节讲六种故障(第 3 节一次讲两种),第 8 节把它们放进一个模拟器里一起玩。
每种故障都按一个格式讲现象 → 后果(带数字)→ 防线 → 在你的项目里怎么制造。
读完以后回到七日冲刺的 D7,翻故障卡、做 applyFills 编程题、写一页总结。
1

故障演练是什么

本节新词故障演练 · 防线 · 注入故障

故障演练就是在出事之前,主动把系统弄坏一次,看它扛不扛得住。

故障演练

生活例子
消防演习:不等真的着火,先拉一次警报,看大家知不知道往哪跑、灭火器在哪、门能不能打开。
算一遍
如果"重复下单"这种问题上线后才发现:一个带单员有 3000 个跟单者,一次重启导致 20 条消息被重新处理,最坏情况下就是 20 × 3000 = 60,000 张重复订单。演练时在本机发现,代价是零。
Java 类比
单元测试里故意 mock 一个会抛异常的依赖,看调用方处理得对不对。故障演练是把这件事放大到整个系统。

防线 和 注入故障

大白话
防线:专门针对某种故障加的保护措施,比如 D6 的幂等编号。
注入故障:在代码里人为制造故障,比如随机丢一条消息、随机让下单超时、在关键两行代码之间杀掉进程。

为什么进组前最值得做这个

背书的人说
"Kafka 可能会重复消费,所以要做幂等。"
做过的人说
"我在提交 offset 之前把进程杀了,重启后第 11–13 条被重新处理。每条要给 3 个跟单者下单,一共 9 张,这 9 张因为编号和第一次一样,全被交易所拒掉了。换成随机编号再试,9 张全成交了。"
区别
后一种说法里有具体的数字、具体的位置和对照实验,一听就知道是真做过的。今天的目标就是让你能这样说话。
常见误区 · 在真环境里演练故障演练只在自己本机或测试环境做:本地的模拟交易所、交易所测试网、本机的 Kafka。绝不在生产环境做,绝不用真钱账户试。进组后想在团队的环境里演练,先问负责人,按团队的流程来。
2

全貌:六个会断的地方

复习的词行情 · 网格 · 扇出 · 带单员 · 跟单者

把 D3 到 D6 做的东西连起来看,一共有六个地方最容易出事。下面的图里用数字标出来了,后面每一节讲一个。

 Exchange WS ──(4)──► [ GridBot ] ──► [ OrderGateway ] ──(3)──► Exchange
                                               │
 leader fill ──► leader.fills ──(1)(2)──► [ Fanout ] ──(5)──► copy.orders ──► [ OrderWorkers ]

 (6) any process may die between "order sent" and "event logged"

图里的名字都是前几天做过的东西:Exchange WS 是交易所推行情的 WebSocket 连接(D3);GridBot 是你的网格机器人(D4);OrderGateway 是负责把订单发给交易所的那一层(D4 项目骨架里的 order 目录)。下面一行是跟单(D6):leader fill 是带单员的一笔成交,先进 leader.fills 这个 topic,Fanout 是扇出服务(D6 的 FanoutService),把它拆成每个跟单者一张订单放进 copy.orders,OrderWorkers 是真正去下单的消费者。括号里的数字对应下表。

编号故障一句话在哪一节
1重复消息同一条成交被处理两遍第 3 节
2rebalance消费者重新分配分区,暂停 + 重投第 3 节
3下单超时不知道单子到底进没进交易所第 4 节
4行情断线 / 跳号本地盘口和真实盘口对不上第 5 节
5热门带单员一个分区被挤爆,跟单者全在排队第 6 节
6进程崩溃做了一半:交易所有单,本地不知道第 7 节
3

重复消息和 rebalance

本节新词重复消息 · rebalance · 去重 · 幂等

重复消息 duplicate

现象
同一条"带单员成交了"被处理了两次。
从哪来
① 消费者处理完、还没提交 offset 就重启了;② 发生了 rebalance(同一个消费者组里有实例加入或退出,分区重新分配),接手的消费者从上一个书签开始读;③ 上游(给你发消息的那一方,这里是发布带单成交的服务)没收到确认,又发了一遍。D6 第 8 节讲过细节。
后果
带单员的仓位统计多算一笔;更严重的是,每个跟单者被多下一张单。一条消息 × 3000 个跟单者 = 3000 张多余的订单。
生活例子
快递员送了两次同一个包裹,你签收了两次,账上扣了两次钱。

两道防线:去重 和 幂等

去重
在自己这边记住处理过的成交 ID(fillId),再来一条同 ID 的就跳过。保护的是你自己的状态(仓位、统计)。
幂等编号
下单时用"成交 ID + 跟单者 ID"做 clientOid,比如成交 ID 是 f1、跟单者是 F17,编号就是 f1-F17。同一条消息不管处理几次,拼出来的编号都一样,交易所会拒绝重复编号(D5 讲过;具体规则和有效时间各家不同)。保护的是交易所那边的订单。
为什么两道都要
扇出服务也可以按 fillId 去重,少扇出一次。但如果它恰好在"订单已经发出"和"记下已处理"之间崩溃,重启后还是会再发一遍。所以下单这一步,最后要靠幂等编号兜底。
关键细节
去重记录要和状态一起保存。如果仓位存在数据库里、去重记录只放在内存的 HashSet 里,一重启 HashSet 就空了,重新投递过来的消息会被当成新消息再算一遍。下面的演示专门演这个。
Java 类比
同一个数据库事务里:INSERT INTO processed_fills(fill_id)(唯一索引)+ UPDATE positions。唯一索引冲突就说明处理过了,整个事务回滚。
一串带重复的成交,最后仓位是多少正确答案:0.7 BTC
算出来的仓位–
正确仓位0.7000
结论–

设定:第 3、6 条是上游重复发送的;最后一次提交 offset 在第 2 条之后,所以重启后从第 3 条开始重新读。正确仓位是不重复的 5 笔相加:0.5 − 0.2 + 0.3 − 0.1 + 0.2 = 0.7。
操作:切换去重方式,再勾掉或勾上"重启",看上面每一行怎么处理、下面的"算出来的仓位"。你应该看到:不去重是 1.5(有重启 2.2);去重记录只放内存,不重启时是 0.7,一重启就变成 1.4;和仓位一起保存,怎么切都是 0.7。这说明去重记录必须和状态一起保存。

在你的项目里怎么制造

重复
生产者发送时按一定概率把同一条发两次(第 9 节有代码)。
rebalance
用同一个 group.id(消费者组的名字,名字相同的实例算一个组,D6 讲过)再启动一个消费者实例,等它分到分区后再 kill -9 掉,观察原来那个消费者的日志和 lag。
4

下单超时:不知道成没成

本节新词超时 · 未知状态 · 按编号查询

超时 和 未知状态

现象
下单请求发出去了,等了规定的时间(比如 3 秒)还没收到回复,这叫超时。这时订单处于未知状态:可能交易所根本没收到,也可能已经收到甚至成交了,只是回复在路上丢了。
生活例子
寄快递后一直查不到物流。可能真寄丢了,也可能只是还没录入。你要是马上再寄一份,收件人可能收到两份。
后果
1000 张单里有 2% 超时,也就是 20 张;假设其中一半其实已经成功了。如果超时就当失败、换个新编号重下,就多出了 10 张重复订单。
Java 类比
HTTP 调用抛了 SocketTimeoutException,不等于对方没执行。和调用支付接口超时是一个道理。

防线:按编号查询

大白话
超时后不要猜,拿同一个 clientOid 去问交易所"有没有这张单"。查到了,就接管它的真实状态:它可能正作为挂单在交易所上等成交,也可能已经成交了;确认查不到,再用同一个编号重新下。万一第一次的请求只是晚到,交易所通常会以"编号重复"拒绝其中一次(D5 讲过),不会出现两张单。
怎么制造
让模拟交易所(D4 的 PaperMatcher,D5 第 9 节改成了存文件的 PaperExchange)按一定概率"订单收下了,但抛出超时"(第 9 节有代码)。
Tips · 查询也会超时按编号查询本身也可能超时或报错。"查询失败"和"查不到"是两回事:查询失败就隔一会儿接着查,这期间这一格不下新单;只有交易所明确回答"没有这张单",才能用原编号重下。
超时以后,两种做法下单 G7-3-B-12 后超时
交易所里这笔单有几张–
结论–

G7-3-B-12 是 D5 讲的可推导编号:7 号网格、第 3 格、买单(B)、第 12 次挂单。
操作:先选做法,再在下拉框里切换交易所的真实情况,看"交易所里这笔单有几张":1 张是对的,2 张就是重复下单。
关键在于:超时那一刻,你不知道下拉框选的是哪一个。做法 A 只在"没收到"时碰巧没事,做法 B 两种情况都对。

录屏:一次下单超时,从发生到防线生效买单 G7-3-B-12 · 交易所其实已经收到
    你的程序以为–
    交易所里这张单–
    这一格一共几张买单–
    这一格能不能下新单–

    录屏:滚到这里会自动一步步播放。每一步变了的数字会被框出来。左边"你的程序以为"和右边"交易所里这张单"一直对不上,直到第 6 步按编号查询之后才对齐;第 7 步是对照,演示如果把超时当失败会发生什么。

    5

    行情断线和跳号

    本节新词行情 · 序列号 · 跳号 · 快照 · 增量 · 行情延迟

    现象

    大白话
    行情就是交易所不停推给你的"订单簿变了、价格变了"的消息。每条消息带一个递增的序列号。连接断了,或者中间丢了一条(序列号从 104 直接跳到 106,叫跳号),你本地的订单簿就和真实的对不上了。D3 讲过。
    生活例子
    看直播球赛,信号卡了 10 秒。恢复后比分牌还是卡住前的数字,你以为还是 0 : 0,其实已经 1 : 0 了。

    后果(算一遍)

    例子
    行情断了 10 秒,这期间真实价格从 60,000 涨到了 60,300,但本地订单簿还停在 60,000。机器人正好在这时铺网格(比如刚重启),按它以为的 60,000,在上方 60,050、60,100 … 60,250 挂了 5 张 0.01 BTC 的卖单。本意是挂单:挂在现价上方,等价格涨上来再卖。
    结果
    真实价格已经在这些价位之上,5 张卖单一挂出去就立刻成交,变成了吃单。成交价是对面挂着的价格(约 60,300),所以价格本身没吃亏。损失在两处:
    ① 手续费:成交额约 5 × 0.01 × 60,300 ≈ 3,015 U,按 D1 的例子(挂单 0.08%、吃单 0.1%)多付约 0.6 U;
    ② 更要命的是位置错了:机器人以为这 5 格还在等上涨,实际一下子全卖了,手里少了 0.05 BTC,后面每一步判断都建立在错的数据上。如果这些卖单带了只做挂单的要求,它们会全部被拒,机器人又会以为自己挂上了。

    防线

    序列号检查
    每条行情都检查序列号是不是上一条 + 1。一旦跳号,整本订单簿丢掉,重新拉一次快照(完整的订单簿),再接着收增量(每条只说"哪一档变成了多少"的变化消息)。
    延迟保护
    记录"最后一次收到行情是什么时候"。行情延迟超过阈值(比如 1 秒)就暂停下单,等行情恢复、重新同步好再继续。
    Java 类比
    TCP 靠序列号发现丢包;这里是应用层自己做一遍。延迟保护就是一个"心跳超时"。
    怎么制造
    最简单:关掉 Wi-Fi 10 秒再打开。或者在收行情的代码里按概率丢掉一条(第 9 节)。
    Tips · 先停手,再修盘口发现跳号或者行情延迟超标,第一件事是暂停这个交易对的下单,第二件事才是重拉快照。顺序反过来,重建盘口的那几百毫秒里,策略还在拿坏数据下单。
    6

    热门带单员:一个分区被挤爆

    复习的词热点分区 · 扇出 · 两段式 · 批量下单

    现象和后果

    现象
    一个热门带单员有 3000 个跟单者。他一成交,就要扇出 3000 张订单。如果这些活都挤在一个分区、一个消费者上做,这个分区的 lag 会猛涨。
    算一遍
    每张单 5 毫秒,3000 张 = 15 秒。排最后的跟单者比第一个晚 15 秒成交。按"价格每秒往不利方向挪 0.02%"的假设,他多付约 15 × 0.02% = 0.3%,一笔 1000 U 的跟单就是 3 U。而且这 15 秒里,同一分区其他带单员的消息也在排队。
    生活例子
    明星开签售会只开一个窗口,队伍排到街上,隔壁窗口的普通业务也被堵住了。

    防线

    两段式扇出
    第一段按带单员分区只做拆分(很快),第二段按跟单者分区,8 个分区并行下单:15 秒 → 约 1.9 秒。
    批量下单
    用"一次提交多张"的接口,减少来回次数。
    怎么制造
    在 followerRepo(D6 代码里存"谁跟了哪个带单员"的那个对象)里给一个带单员塞 3000 个跟单者,发一条成交,看第二段每个分区的 lag 和最后一张单的时间。
    7

    进程崩溃:写了一半

    本节新词崩溃 · 事件日志 · 孤儿单 · 重启对账

    现象

    大白话
    程序在两步操作之间突然死掉(崩溃),比如被 kill -9、机器断电。最危险的位置是"订单已经发给交易所"和"本地记下这件事"之间:交易所有这张单,本地却不知道。这种单叫孤儿单。
    生活例子
    你刷卡付了钱,正要在记账本上记一笔,手机没电关机了。第二天看账本,这笔钱对不上。

    后果(算一遍)

    例子
    网格挂着 10 张单,崩溃时有 1 张"已经发出、还没记下"。重启后机器人以为那一格是空的,又挂了一张:同一格有了两张买单。等它们都成交,网格就多买了 0.01 BTC,之后每一步都会跟着错。

    防线

    先写日志
    事件日志是一个只追加、不修改的文件,每件事先写进去,再改内存里的状态。重启时从头回放一遍,就能恢复崩溃前的状态。D5 讲过。
    重启对账
    日志只能保证"记下来的事不丢",没来得及记的事还是会丢。所以重启后要和交易所核对一遍:交易所有、本地没有的挂单(孤儿单),先按 clientOid 前缀确认是不是自己的(比如编号以 G7- 开头,就是 7 号网格发的;D5 的 ClientOid.belongsTo 就是干这个的),再决定认领还是撤销。
    Java 类比
    数据库的预写日志(WAL):先写日志再改数据页,崩溃后靠日志恢复。
    怎么制造
    在"下单"和"写日志"之间加一行 Runtime.getRuntime().halt(1)(第 9 节有代码),它比 System.exit 更接近 kill -9:不执行任何收尾代码。
    8

    故障模拟器:选一种故障,开关防线

    把第 3–7 节的六种故障放在一起。选一种,勾上或去掉对应的防线,点"注入故障"看后果。

    故障模拟器数字是按每节里的例子设定的

    操作:先选一种故障,读"场景",点"注入故障"看结果(默认所有防线都开着);然后去掉一道防线再点一次,对比数字怎么变。结论标签有三种:扛住了、部分扛住、出事了。

    9

    在你的 Java 项目里注入故障

    做法很简单:加一个开关类 Chaos(英文"混乱"),平时全部关掉;演练时打开某一个开关,看系统的反应。下面的类名和方法名沿用 D3–D6 笔记里的代码。

    // 故障开关:演练时打开,平时全部为 0 / false
    import java.util.Random;
    
    public final class Chaos {
        public static volatile double DUPLICATE_RATE = 0.0;       // 重复发送消息的概率
        public static volatile double DROP_RATE = 0.0;            // 丢掉一条行情的概率
        public static volatile double LOST_ACK_RATE = 0.0;        // 下单成功但回复丢了的概率
        public static volatile boolean CRASH_AFTER_PLACE = false; // 下单之后、写日志之前直接死掉
    
        private static final Random R = new Random();
        public static boolean hit(double rate) { return R.nextDouble() < rate; }
    }

    重复发送(对应第 3 节)

    // rec 是要发的那条 Kafka 消息,producer 是 D6 里的 KafkaProducer
    void publish(ProducerRecord<String, String> rec) {
        producer.send(rec);
        if (Chaos.hit(Chaos.DUPLICATE_RATE)) producer.send(rec);   // 同一条再发一次
    }

    下单超时(对应第 4 节)

    // D5 第 9 节的 PaperExchange(存文件的模拟交易所):订单收下了,但故意抛超时
    public Order place(Order o) throws TimeoutException {
        if (open.containsKey(o.clientOid()))                 // 编号已存在:像真交易所一样拒绝
            throw new IllegalStateException("DUPLICATE_CLIENT_OID");
        open.put(o.clientOid(), o);                          // 订单已经进了"交易所"
        save();                                              // 写 exchange.db
        if (Chaos.hit(Chaos.LOST_ACK_RATE))
            throw new TimeoutException("模拟:成功了,但回复丢了");
        return o;
    }
    
    // 调用方:就是 D5 第 4 节的 placeSafely,超时后用同一个编号先查
    public Order placeSafely(Order o) throws Exception {
        try {
            return exchange.place(o);
        } catch (TimeoutException e) {
            Optional<Order> found = exchange.query(o.clientOid());
            if (found.isPresent()) return found.get();       // 已经有了:接管真实状态
            return exchange.place(o);                        // 确认没有:用同一个编号重下
        }
    }

    丢行情(对应第 5 节)

    // D3 第 10 节处理增量的方法,在最前面加一行
    void onDelta(Delta d) {
        if (Chaos.hit(Chaos.DROP_RATE)) return;             // 假装这条在网络上丢了
        // ……下面是 D3 原来的逻辑:先检查序列号,跳号就丢掉订单簿、重新拉快照
    }

    写了一半就死(对应第 7 节)

    Order placed = exchange.place(o);
    if (Chaos.CRASH_AFTER_PLACE) Runtime.getRuntime().halt(1);  // 不跑任何收尾代码,接近 kill -9
    local.record("PLACED " + placed.clientOid());              // D5 的 LocalBook:先写日志,再改内存

    演练顺序建议

    每次只开一个
    一次只打开一个开关,否则出了问题分不清是谁造成的。
    先预测再动手
    动手前先写下"我预计会看到什么",跑完对比。预测错了的地方,就是你真正学到东西的地方。
    留下证据
    每次演练记下:开了哪个开关、看到了什么日志、数字是多少。这些就是第 11 节一页总结的素材。
    10

    监控要看什么

    本节新词监控 · 告警 · 阈值 · 行情延迟 · 下单失败率 · 对账差异数

    演练是你主动去找问题,监控是让问题自己冒出来。系统持续记录几个关键数字,超过事先定好的阈值就发告警(发消息、打电话给值班的人)。

    四个最重要的数字

    lag
    Kafka 里还没处理的消息有多少条。涨了说明跟不上,跟单者在排队。
    行情延迟
    现在距离最后一次收到行情过了多久。大了说明行情断了或者卡了,策略在用旧价格做决定。
    下单失败率
    下单请求里被拒绝或出错的比例。突然升高,可能是额度规则变了、交易所限流了,或者你的代码出了 bug。
    对账差异数
    定时拿本地记录和交易所的真实记录比,对不上的有几条。正常应该永远是 0,出现一条就要查。
    Java 类比
    和你熟悉的接口 QPS、错误率、P99 延迟监控是一回事,只是换成了交易系统关心的数字。
    拖一拖:什么时候该告警示例阈值,真实数值以你们组的配置为准

    操作:拖每一行的滑块,看右边的标签在哪个数值从"正常"变成"关注"、再变成"告警"。注意对账差异数:只要不是 0,就直接告警。

    11

    一页总结怎么写

    七天的成果最后落到一页纸上:这个系统会在哪 5 个地方出事,每个都写清现象、后果、防线和你亲手验证的证据。进组第一周,它就是你跟同事聊天的底稿。

    故障名称:
      现象:    一句话,别人能听懂发生了什么
      后果:    带数字,比如"多下了 9 张单""多付了 0.6 U 手续费"
      防线:    用了什么机制,写在代码的哪里
      我的验证:开了哪个开关,看到了什么,对照实验是什么
      还没想明白:进组后要问的问题

    写好的示例

    故障名称
    扇出服务重复处理带单成交
    现象
    扇出服务处理完消息、还没提交 offset 就重启,重启后把最后几条再处理了一遍。
    后果
    每条消息扇出给 3 个跟单者。每 5 条提交一次,在第 13 条后崩溃,第 11–13 条被重新处理,重新发出 9 张订单。
    防线
    clientOid = 带单员成交 ID + "-" + 跟单者 ID(写在扇出服务 FanoutService 里拼编号的那一行);下单消费者按 clientOid 拒绝重复。
    我的验证
    在扇出服务的 producer.flush()(确认消息都发出去了)和 commitSync()(提交 offset,也就是夹书签)之间加 halt(1),重启后日志里 9 行"重复订单被拒绝"。把 clientOid 换成 UUID 做对照,9 张全部成交。
    还没想明白
    真实交易所对重复 clientOid 的去重时间窗口有多长?我们组有没有自己再做一层去重?

    写的时候直接写在 七日冲刺 D7 的总结框里,会自动保存在本机浏览器。

    12

    词典、自测和七天收尾

    今天用到的词都在这里,包括前几天学过、今天又用到的。

    自测 7 题

    七天结束

    回到 七日冲刺 D7:翻一遍六张故障卡,做 applyFills 编程题,然后把一页总结写完。最后看 进组以后,那里有第一、二周的安排。

    国庆时间还够的话,接着做三天加练:D8 交易所内部怎么跑现货和合约,D9 量化策略怎么判断好坏,D10 在测试网上亲手操作一遍。

    进组第一周怎么说话,三个例子:

    像背书"Kafka 是至少一次语义,所以消费端要做幂等。"

    像干过"咱们跟单下单用的 clientOid 是怎么拼的?我本地试过用随机编号,重启后会重复下单。"

    像背书"要监控消费延迟,防止消息堆积。"

    像干过"头部带单员成交的时候,copy.orders 各个分区的 lag 一般会冲到多少?有没有告警阈值?"

    像背书"下单超时要保证幂等性。"

    像干过"下单超时以后,我们是先按 clientOid 查一次,还是直接等私有推送的订单回报?"