程序一定会挂
机器人是 7×24 小时跑的,它一定会在某个时刻突然死掉。问题不是"会不会挂",而是"挂了以后重启,能不能接着干,而且不出错"。
崩溃 和 kill -9
finally 块和关闭钩子都不会执行。内存状态 和 持久化
持久化:把状态存到进程死了也不会丢的地方,比如磁盘文件、数据库。
事件日志:先记账,再改内存
把"发生了什么事"一行一行写进文件,只往后加、从不修改。重启时从第一行读到最后一行,每一行重新执行一遍,内存状态就回来了。
事件日志 和 只追加
PLACED A BUY 59000(挂了 A)PLACED B SELL 61000(挂了 B)FILLED A(A 成交了)PLACED C SELL 60000(挂了 C)从头读一遍:{A} → {A, B} → {B} → {B, C}。结论:现在挂着 B 和 C。
为了好读,这里省略了格子和数量。下面代码里的真实格式是"PLACED 编号 方向 格子 价格 数量",比如
PLACED G1-2-B-1 BUY 2 59000 0.01。其中 G1-2-B-1 是我们自己给这张单起的编号,叫 clientOid(第 5 节细讲)。回放 和 预写日志 replay / WAL
apply 方法。顺序一反,崩在两步中间的那件事就永远丢了。交易所上实际发生的事
日志文件(磁盘)
内存状态:本地以为在挂的单
录屏:滚到这里会自动播一遍。前 5 步是正确顺序(崩溃后回放,和交易所一致),后 4 步是错误顺序(崩溃后冒出一张孤儿单)。随时可以暂停,或用"上一步 / 下一步"来回看。看完想自己动手,就点上面那几排按钮:一点它们,录屏就停下,交给你操作。
自己动手:先用默认的"先写日志",点两次"处理下一个事件",再点"中途崩溃",然后"重启并回放":右边的内存和左边交易所的实际情况一致。再切到"先改内存",在第 1 步或第 3 步点"中途崩溃"再回放:会看到回放后对不上,分别出现孤儿单和"以为还在挂"的单。这说明顺序决定了崩溃后能不能恢复。
代码分两个类。EventLog 只管日志文件本身:追加一行、从头读一遍。LocalBook 是"本地以为在挂的单",它用 EventLog 来记账。
关键在于,正常运行和重启回放用的是同一个 apply 方法,所以两种情况下得到的状态一定一样。
另外一个细节:flush() 只是把数据交给操作系统,操作系统可能还没真正写到磁盘;进程被 kill 不影响,但整台机器断电就可能丢最后几行。要扛断电,得再调 FileChannel.force(true) 强制写盘,代价是每次写都慢很多。练习里用 flush 就够了。
public class EventLog implements AutoCloseable {
private final Path file;
private final BufferedWriter out;
public EventLog(Path file) throws IOException {
this.file = file;
this.out = Files.newBufferedWriter(file, StandardCharsets.UTF_8,
StandardOpenOption.CREATE, StandardOpenOption.APPEND); // 只追加
}
/** 写一行,flush 之后才返回 */
public void append(String line) throws IOException {
out.write(line);
out.newLine();
out.flush(); // 交给操作系统(见上面的说明)
}
/** 启动时从头读一遍,每一行交给 handler */
public void replay(Consumer<String> handler) throws IOException {
if (!Files.exists(file)) return;
try (Stream<String> lines = Files.lines(file, StandardCharsets.UTF_8)) {
lines.forEach(handler);
}
}
@Override public void close() throws IOException { out.close(); }
}
public class LocalBook {
private final Map<String, Order> open = new LinkedHashMap<>(); // 本地以为在挂的单
private final EventLog log;
public LocalBook(EventLog log) throws IOException {
this.log = log;
log.replay(this::apply); // 重启:回放
}
/** 正常运行:先写日志,再改内存 */
public void record(String line) throws IOException {
log.append(line);
apply(line);
}
/** 挂单成功后记一行,格式:PLACED 编号 方向 格子 价格 数量 */
public void recordPlaced(Order o) throws IOException {
record("PLACED " + o.clientOid() + " " + o.side() + " " + o.level()
+ " " + o.price().toPlainString() + " " + o.qty().toPlainString());
}
/** 回放和正常运行共用这一个方法 */
private void apply(String line) {
String[] a = line.split(" ");
switch (a[0]) {
// PLACED 编号 方向 格子 价格 数量
// Order 就是 D4 的 record Order(clientOid, side, price, qty, level),注意参数顺序
case "PLACED" -> open.put(a[1], new Order(a[1], Side.valueOf(a[2]),
new BigDecimal(a[4]), new BigDecimal(a[5]), Integer.parseInt(a[3])));
case "FILLED", "CANCELED" -> open.remove(a[1]);
default -> throw new IllegalStateException("看不懂的日志: " + line);
}
}
public Collection<Order> openOrders() { return open.values(); }
}
快照:日志太长怎么办
快照 snapshot
下单超时:你不知道发生了什么
下单请求发出去,等了 3 秒没有回音。这张单到底下成功没有?你不知道。这是交易系统里最常见、也最容易出错的情况。
请求、响应 和 超时
超时只告诉你"没等到",不告诉你是哪一种。
HttpTimeoutException、SocketTimeoutException:异常只说明你这边没等到,不说明对方做没做。和你做支付时"调用支付渠道超时"是一回事。 place(clientOid = G1-3-S-7) --> TIMEOUT ??
|
+-- [X] retry with a NEW clientOid -> may create a 2nd order
|
+-- [OK] query by the SAME clientOid
+-- found -> adopt its real status (OPEN / FILLED)
+-- not found -> re-place with the SAME clientOid
(a late duplicate gets rejected)
错误做法:换一个新编号重试。正确做法:用同一个编号去查,查到就接管它的真实状态,查不到再用同一个编号重新下。图里的 G1-3-S-7 是我们自己给这张单起的编号,代码里叫 clientOid(网格 G1、第 3 格、卖单、第 7 次,第 5 节细讲)。
怎么玩:把上面四个按钮依次点一遍,对比左右两栏。左边"换新编号"只在 ① 碰巧没事,② ③ ④ 都出错;右边"同一个编号先查"四种情况都对。超时的时候你并不知道是哪一种,所以只能选"四种都对"的那个做法。
写代码前,先把"交易所"抽象成一个接口,叫 。今天它的实现是模拟交易所(第 9 节),进组后换成真实交易所的接口,上层代码不用改。
接口里只有四件事:下单、按 clientOid 查单、撤单(取消一张还没成交的挂单)、查所有还在挂的单。查回来如果发现已经成交,要按成交处理:记日志,再挂反向单(D4 讲过:一格成交后,在相邻格挂一张反方向的单)。这和第 7 节的对账(拿本地记录和交易所逐条比)用的是同一套规则。
/** 交易所查回来的结果:订单本身 + 它在交易所的状态 */
public record RemoteOrder(Order order, String status) {} // OPEN / PARTIALLY_FILLED / FILLED / CANCELED
public interface Exchange {
Order place(Order o) throws TimeoutException; // 下单;没等到响应就抛 TimeoutException(java.util.concurrent 里的)
Optional<RemoteOrder> query(String clientOid); // 按编号查;交易所没有这张单就返回 Optional.empty()
void cancel(String clientOid); // 撤单
List<RemoteOrder> openOrders(); // 交易所上所有还在挂的单
}
/** 下单,遇到超时也不会下重复单。返回这张单在交易所的真实状态 */
public String placeSafely(Order o) throws Exception {
try {
exchange.place(o); // 正常情况:挂上了
return "OPEN";
} catch (TimeoutException e) {
Optional<RemoteOrder> found = exchange.query(o.clientOid()); // 用同一个编号去查
if (found.isPresent()) return found.get().status(); // 其实已经下成功了:接管它的真实状态
exchange.place(o); // 查不到:用同一个编号重下
return "OPEN";
// 如果第一次的请求只是晚到,交易所会以"编号重复"拒绝其中一次,不会出现两张单
}
}
// 调用方拿到 "FILLED" 时,要按成交处理:记日志 + 挂反向单(和第 7 节对账里的做法一样)
自己起订单编号:clientOid
orderId 和 clientOid
clientOid:你自己生成、跟着下单请求一起发过去的编号。交易所会记下来,之后可以拿它来查询、撤单。
可推导的编号
G1-3-S-7 = 网格 G1、第 3 格、卖单、这一格第 7 次挂单。假设进程死在"发出第 7 次的请求"和"写日志"之间。重启回放,日志里第 3 格最后一条是"第 6 次卖单已成交"。
那么如果还有下一张,它一定是第 7 次,编号一定是
G1-3-S-7。拿这个编号去交易所查一下,就知道它到底挂上没有。3f2a9c1e-…),没法从已知信息算出来。同样死在"发请求"和"写日志"之间,这串随机字符没被记下来,你就再也不知道它是什么,也就没法去查。public final class ClientOid {
/** 网格-格子-方向-次数,例如 G1-3-S-7 */
public static String of(String gridId, int level, Side side, int cycle) {
return gridId + "-" + level + "-" + side.name().charAt(0) + "-" + cycle;
}
/** 判断一个编号是不是属于这个网格(第 7 节对账时用) */
public static boolean belongsTo(String clientOid, String gridId) {
return clientOid.startsWith(gridId + "-");
}
}
各家交易所对 clientOid 的长度和可用字符有限制,常见是只允许字母、数字、横线,长度在几十个字符以内,具体看你们对接的交易所文档。
幂等:同一件事做多次,结果和做一次一样
幂等 和 去重
每次换新编号 → 3 张单,这一格买了 3 倍,占用 1,770 U。
每次用同一个编号 → 1 张成功,另外 2 次被拒"编号重复",占用 590 U。收到"编号重复"时,你应该理解成"其实已经下过了"。
怎么玩:先拖"发了几次",再在两种做法之间切换。换新编号时,张数和占用资金跟着次数一起涨;用同一个编号,永远只有 1 张、590 U,多出来的请求都被拒绝。
重启对账:以交易所为准
重启回放之后,本地知道了"我以为在挂的单"。但这只是本地的记录,真实情况要以交易所为准。把两边逐张对一遍,差异分类处理,这就是对账。
对账 和 以交易所为准
孤儿单 和 撤单
G1- 开头),别把别的网格或别人的单撤了。| 本地 | 交易所 | 处理 | 原因 |
|---|---|---|---|
| 在挂 | 挂着 / 部分成交 | 保留 | 两边一致 |
| 在挂 | 已全部成交 | 标记成交 | 记一条 FILLED,并走网格下一步:挂反向单 |
| 在挂 | 已撤 | 重新下(新编号) | 那张单已经结束了,这一格现在是空的,按原计划补一张,次数加 1 |
| 在挂 | 查不到 | 重新下(原编号) | 说明当初根本没挂上;用原编号重下,和第 4 节一样,万一原请求只是晚到,交易所会拒绝重复 |
| 不知道 | 挂着,编号属于本网格 | 撤掉孤儿单 | 本地不认识它,留着会让状态跑偏 |
| 不知道 | 挂着,编号不属于本网格 | 不动 | 不是你的单 |
怎么玩:默认三张单在交易所分别是"挂着""已全部成交""查不到",正好对应三种处理。改一改第三列的下拉框,看最后一列怎么变;再取消勾选下面两个复选框,看孤儿单和别的网格的单那两行消失。
下面的 Reconciler(对账器)就是上面那张表的代码版。它用到第 4 节的 Exchange 接口、第 2 节的 LocalBook、D4 的 和 Fill。
public class Reconciler {
private final GridStrategy strategy;
private final Exchange exchange;
private final LocalBook local;
public Reconciler(GridStrategy strategy, Exchange exchange, LocalBook local) {
this.strategy = strategy; this.exchange = exchange; this.local = local;
}
public void reconcile(String gridId) throws Exception {
// 1. 本地以为在挂的,逐张去交易所查(copyOf:边遍历边改会抛异常,先拷一份)
for (Order o : List.copyOf(local.openOrders())) {
String s = exchange.query(o.clientOid()).map(RemoteOrder::status).orElse("NOT_FOUND");
switch (s) {
case "OPEN", "PARTIALLY_FILLED" -> { } // 保留
case "FILLED" -> {
local.record("FILLED " + o.clientOid()); // 先记日志
Optional<Order> next = strategy.nextOrderOnFill(toFill(o));
if (next.isPresent()) placeAndRecord(next.get()); // 挂反向单
}
case "CANCELED" -> { // 被撤了:这一格空了
local.record("CANCELED " + o.clientOid());
placeAndRecord(withNextCycle(o)); // 新编号,次数 +1
}
default -> { // NOT_FOUND:当初没挂上
local.record("CANCELED " + o.clientOid()); // 先从本地记录里去掉
placeAndRecord(o); // 用原编号重下(同第 4 节)
}
}
}
// 2. 交易所上在挂、本地不知道的:确认是本网格的,就撤掉
Set<String> known = local.openOrders().stream()
.map(Order::clientOid).collect(Collectors.toSet());
for (RemoteOrder r : exchange.openOrders()) {
String id = r.order().clientOid();
if (!known.contains(id) && ClientOid.belongsTo(id, gridId)) {
exchange.cancel(id); // 孤儿单
}
}
}
/** 下单成功后记日志(今天先用这个简单顺序,第 8 节再讲"意图") */
private void placeAndRecord(Order o) throws Exception {
exchange.place(o);
local.recordPlaced(o);
}
/** 把一张挂单包装成 D4 的"全部成交"回报,交给策略算下一张 */
private static Fill toFill(Order o) {
return new Fill(o.clientOid(), o.side(), o.level(), o.price(), o.qty(), "FILLED");
}
/** 同一格、同一方向,次数 +1:G1-3-S-7 → G1-3-S-8 */
private static Order withNextCycle(Order o) {
String id = o.clientOid();
int cut = id.lastIndexOf('-');
int n = Integer.parseInt(id.substring(cut + 1)) + 1;
return new Order(id.substring(0, cut + 1) + n, o.side(), o.price(), o.qty(), o.level());
}
}
为什么"已撤"用新编号、"查不到"用原编号:已撤的那张单已经结束了,再挂是一件新的事,换编号;查不到说明当初根本没挂上,还是同一件事,用原编号,万一原请求只是晚到,交易所会拒绝重复。练习里 withNextCycle 直接在编号末尾加 1;真实项目里次数最好统一由 GridStrategy 管,避免两处各算各的撞号。
重启的完整流程
start | v [1] load snapshot (if any) <- 第 3 节,今天可以跳过 v [2] replay EventLog <- 第 2 节,得到本地以为在挂的单 v [3] query exchange by clientOid <- 第 5 节 v [4] reconcile <- 第 7 节 v [5] subscribe market data, start loop
对账做完之前,不要开始处理行情。否则你是拿着一份不准的状态在做决策。
下单意图:把孤儿单消灭在源头
G1-3-S-7 去交易所查一下。孤儿单基本就不会出现了。这也是为什么编号要可推导:意图里记的就是这个编号。下面的 main 把前面几节串起来。里面有两个名字要先认识: 是第 9 节写的、会把挂单存进文件的模拟交易所,用来代替 D4 里做模拟撮合的 PaperMatcher; 是 D4 的事件循环:一个线程按顺序取行情和成交回报来处理。
public static void main(String[] args) throws Exception {
EventLog log = new EventLog(Path.of("grid-G1.log"));
LocalBook local = new LocalBook(log); // [2] 构造时就回放日志
Exchange exchange = new PaperExchange(Path.of("exchange.db")); // 见第 9 节
GridStrategy strategy = GridStrategy.restore("G1", local); // 按本地挂单恢复每一格状态(今天要加的方法,见下)
new Reconciler(strategy, exchange, local).reconcile("G1"); // [3][4] 先对账
BotLoop loop = new BotLoop(exchange, strategy, local); // [5] 再开始干活(构造参数比 D4 多了,见下)
new Thread(loop, "bot-loop").start();
PublicWsClient.connect(price -> loop.submit(new Trade(price)));
new CountDownLatch(1).await(); // 让 main 线程一直等着,不退出
}
和 D4 相比,今天要改三处:① GridStrategy 加一个静态方法 restore(gridId, local):按本地在挂的单,把对应格子设成买单或卖单,并把每格的次数设成编号里最大的那个数;② BotLoop 的构造参数从 D4 的 (PaperMatcher, GridStrategy) 改成 (Exchange, GridStrategy, LocalBook),挂新单时走"下单 + recordPlaced 记日志",收到成交时先 record("FILLED …");③ PaperMatcher 换成第 9 节那个会存文件的 PaperExchange。
今天的完成标准和怎么自测
完成标准只有一句话:运行中直接 kill 掉进程,重启后状态和挂单能恢复,而且没有重复下单。
先解决一个问题:模拟交易所也会跟着死
PaperMatcher 是个内存里的 Map。进程一死,"交易所"也跟着没了,你就没法练对账:重启后问它,它什么都不记得。Exchange 接口,内部还是 D4 那套"价格碰到就成交"的逻辑,但把挂单单独存一个文件(比如 exchange.db,名字随意,其实就是个文本文件),每次变化都整体写一遍,启动时读回来。这样它就像一个"不会跟着你一起死的交易所"。注入故障 fault injection
// 自测 1:在 PaperExchange 里,随机"挂上了但不告诉你",逼出超时逻辑
@Override
public Order place(Order o) throws TimeoutException {
open.put(o.clientOid(), o);
save(); // 写 exchange.db
if (random.nextInt(10) == 0) throw new TimeoutException("假装超时"); // 10% 概率
return o;
}
// 自测 2:在"下单"和"记日志"之间直接死掉,制造一张孤儿单(改 placeAndRecord)
exchange.place(o);
if (System.getenv("CRASH_AFTER_PLACE") != null) Runtime.getRuntime().halt(1); // 立刻退出,不跑任何收尾,效果等同 kill -9
local.recordPlaced(o);
CRASH_AFTER_PLACE 是一个环境变量开关:启动时设了它(CRASH_AFTER_PLACE=1 java …)就会在那一行死掉,不设就正常运行,不用改代码。
- 正常跑一会儿,
kill -9 <进程号>,重启:本地挂单和exchange.db里的一致,每格最多一张 - 打开自测 1 跑 10 分钟:超时后没有出现同一格两张单
- 用
CRASH_AFTER_PLACE=1启动一次,进程死后去掉这个变量重启:对账撤掉了那张孤儿单 - 停机期间手动改
exchange.db,把一张单改成已成交,重启:对账把它标记成交并挂了反向单
词典和自测
今天用到的词,包括前几天学过的,这里都能查到。
自测 7 题
回到 七日冲刺 D5:先看"下单超时怎么办"那张图,对照第 4 节的演示;再看"重启对账"那张图,对照第 7 节的表格;最后做 reconcile 编程题,它就是第 7 节那段 Java 代码的 JavaScript 版本。