今天要做什么:跟单从头到尾
今天把 D1 第 11 节讲的跟单真正用代码跑起来。先把业务流程完整走一遍,后面每一节都会对应到其中某一步。
录屏:滚到这里会自动一步步播放,每步停 4 秒左右,高亮的是当前这一步。想细看就点"暂停",用"上一步 / 下一步"来回翻。
第 2–6 节讲 Kafka 怎么把"成交 → 找人 → 下单"这几步连起来;第 7 节讲"每人下多少";第 8 节讲下单重试时最容易踩的坑。
为什么要用消息队列
消息队列是一个"先放着、慢慢处理"的中转站。发的一方把消息放进去就走,处理的一方按自己的速度一条条取出来处理。
消息队列 message queue
BlockingQueue 一样:一边 put(),一边 take()。区别是 Kafka 是独立部署的服务,消息存在磁盘上,程序重启了消息还在,多台机器都能来读。生产者 和 消费者 producer / consumer
消费者:从队列里取消息来处理的程序,这里是跟单服务。
Kafka 的 topic、分区和 key
Kafka 是最常用的消息队列之一,你们组用得很多。先记住三个词:topic 是"哪一类消息",分区是"这类消息分成几条队来排",key 决定"这条消息排哪一队"。
topic 主题
Map<String, Queue> 里的那个 String。分区 partition
List<BlockingQueue> partitions,每个队列配一个线程。key 和 同分区有序
int p = toPositive(hash(key)) % partitions; Kafka 默认用一种叫 murmur2 的哈希算法,思路和 HashMap 选桶一样。 topic: leader.fills (4 partitions)
producer ── send(key="L1") ──► hash("L1") % 4 = 2
│
P0 [ m3 ][ m7 ] │
P1 [ m1 ][ m5 ][ m9 ] ▼
P2 [ L1-a ][ L1-b ][ L1-c ] ← L1 的消息全在这里,按顺序
P3 [ m2 ][ m8 ]
图里 P0–P3 是 4 个分区;m1、m2… 是别的带单员的消息;L1-a、L1-b、L1-c 是 L1 的三条消息,全在 P2 里按顺序排着。
消费者组、offset 和 lag
消费者组 consumer group
offset 和 提交 offset / commit
但如果它处理到第 79 条才挂,书签还停在 60,那 60–79 这 20 条会被再处理一遍。第 8 节专门讲这个。
consumer.commitSync() 就是夹书签。enable.auto.commit)。它只看时间,不看你处理完没有:书签可能已经挪到第 80 条,你才处理到第 70 条,这时程序挂了,第 70–79 条就被跳过了,也就是丢消息。练习里一律关掉自动提交,处理完再手动 commitSync():宁可重复(第 8 节有防线),不能丢。lag 消费延迟、堆积量
size()。怎么玩:默认只有 2 个消费者,每秒处理 800 条、进来 1200 条,60 秒后堆了 24,000 条。把"消费者数"拉到 3,堆积消失;再拉到 6,看"真正在干活的消费者"停在 4,因为分区只有 4 个,多出来的人闲着。
分区键怎么选
分区键就是发消息时拿哪个字段当 key。选哪个字段,决定了两件事:谁和谁的消息一定按顺序处理,以及流量会不会都挤进同一个分区。
| 用什么当 key | 保证什么顺序 | 问题 |
|---|---|---|
| 带单员 ID | 同一个带单员的成交按顺序处理(开仓一定在平仓之前) | 热门带单员流量大,他所在的分区会堆积 |
| 交易对 | 同一个交易对的消息按顺序 | 跟单用不上这个顺序;而且 BTC 流量最大,同样会挤 |
| 跟单者 ID | 同一个跟单者的订单按顺序 | 流量均匀。但前提是消息已经拆成了"每个跟单者一条"(第 6 节) |
热点分区 hot partition
HashMap 里大量 key 撞到同一个桶,那个桶的链表特别长。怎么玩:在下拉框里依次切换三种 key,每种看几秒。条越长 = 这个分区堆的越多;变红 = 堆积超过 40 条。每个分区的消费者每一拍处理 4 条,总共进来 12 条,平均分的话完全够用。按带单员或交易对时,会有一条一直变长变红(热点分区);按跟单者时,四条都很短。
两段式扇出
扇出 fan-out
for (Follower f : followers) placeOrder(f);两段式
第一段只负责把一条带单成交拆成 N 条跟单订单(只是算数和发消息,很快);
第二段按跟单者 ID 分区,多个消费者并行真正去下单(慢的活摊开做)。
producer.send()(很快);第二段是多个消费者实例,各管几个分区。 leader.fills ┌───────────────┐ copy.orders ┌──────────────┐
key = leaderId ───────────► │ FanoutService │ ───► key = followerId ───────► │ OrderWorkers │
└───────────────┘ └──────────────┘
ordered per leader 1 msg ──► N msgs spread evenly place orders
(fast: only math) ordered per follower in parallel
图里的英文:FanoutService = 扇出服务(第一段,只拆不下单);OrderWorkers = 下单消费者(第二段,真正去下单);leaderId = 带单员 ID;followerId = 跟单者 ID。"ordered per leader"= 同一个带单员的消息按顺序,"spread evenly"= 流量均匀摊开。
怎么玩:把"第二段分区数"从 1 拖到 16,看最后一人等待的时间和多付的钱怎么下降;分区数是 1 时,两段式和单线程一样慢。"多付"按一笔 1000 U 的跟单估算:等得越久,价格挪得越远。0.02%/秒只是为了演示的假设,真实行情波动时可能大得多。
每个人下多少
第一段拆消息时,要给每个跟单者算出"下多少个"。规则只有四条,但每一条都跟钱有关。
四条规则
例:带单员本金 10000 U 买了 1 个 BTC;你本金 5000 U、倍数 1 → 1 × 0.5 × 1 = 0.5 个。
例:固定跟 200 U,价格 60000 → 200 ÷ 60000 = 0.003333… 个。
BigDecimal,见下面的代码。// LeaderFill(带单员成交)和 Follower(跟单者)的定义在第 10 节第 3 步
public class CopySizer {
/** 按比例或固定金额算出原始数量,再向下取整到 lotSize */
public static BigDecimal size(LeaderFill fill, Follower f, BigDecimal lot) {
BigDecimal raw = f.mode() == Mode.RATIO
? fill.qty().multiply(f.equity()).multiply(f.value()) // 带单数量 × 跟单者本金 × 倍数
.divide(fill.leaderEquity(), 12, RoundingMode.DOWN) // ÷ 带单员本金
: f.value().divide(fill.price(), 12, RoundingMode.DOWN); // 固定金额 ÷ 价格
return raw.divide(lot, 0, RoundingMode.DOWN).multiply(lot); // 除以 lot 取整数份,再乘回去 = 向下取整
}
}
额度 = 这个人现在最多能下多少 U 的单(合约里已经算上了)。默认 F3 太小被跳过、F4 额度不够被跳过。试试把 F3 的固定金额改成 60(刚好够 0.001 个),或者把 F4 的额度改成 20000,看它们变成可以下单。
重复消息从哪来,怎么防
Kafka 保证消息不丢,但不保证只处理一次。同一条"带单员成交了"被处理两遍(这就叫重复消息),跟单者就会被下两次单,仓位翻倍。这是跟单系统里最危险的 bug 之一。
rebalance 再平衡
consumer.subscribe(topics, new ConsumerRebalanceListener(){...}),回调里能看到自己丢了哪些分区、拿到哪些分区。至少一次 at-least-once
幂等 idempotent
clientOid 客户端订单号
clientOid = 带单员成交ID + "-" + 跟单者ID,比如 = 带单员成交 LF-1001 给跟单者 F7 下的单。同一笔带单成交、同一个跟单者,不管重试多少次,算出来的编号都一样,所以不会重复下单。UUID.randomUUID() 当编号:每次重试编号都不同,交易所就当成新订单,照单全收。enable.idempotence 只管"生产者自己重发"这一侧,管不了消费者重复处理。所以别指望"每条只处理一次",把最后一道防线放在下单编号上:clientOid = 带单员成交 ID + 跟单者 ID,重复多少次都是同一个编号。怎么玩:默认每 5 条提交一次、第 13 条后崩溃,黄色的第 11–13 条会被再处理一遍。把"每处理几条提交一次"拖到 1,重复就没了,但每条都要和 Kafka 来回一次,更慢;再把它拖回 5,勾上"用幂等编号",重复处理还在,但"交易所实际多成交"变成 0。这说明重复很难完全避免,真正的防线在下单编号。
公平性:谁先下单
3000 个跟单者,总有人先下单、有人后下单。价格一直在变,先下单的人价格更好。让谁排第一,就是公平性问题。
公平性
常见做法
fill #1: F1 F2 F3 F4 F5 fill #2: F2 F3 F4 F5 F1 ← 轮转:每次换一个人排第一 fill #3: F3 F4 F5 F1 F2
动手:在本机跑起来
目标:发一条"带单员成交了",看到 3 个跟单者各自收到一张订单;中途杀掉程序再启动,看到重复的订单被拒绝。
第 1 步:起 Kafka
# 启动 Kafka(单机模式,端口 9092) docker run -d --name kafka -p 9092:9092 apache/kafka:3.8.0 # 建两个 topic,各 4 个分区 docker exec kafka /opt/kafka/bin/kafka-topics.sh --create --topic leader.fills \ --partitions 4 --bootstrap-server localhost:9092 docker exec kafka /opt/kafka/bin/kafka-topics.sh --create --topic copy.orders \ --partitions 4 --bootstrap-server localhost:9092
第 2 步:加依赖
KafkaProducer(生产者)和 KafkaConsumer(消费者)。JSON 解析用你熟悉的 Jackson 就行,下面的代码里省略了。<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.8.0</version> </dependency>
第 3 步:先定义几个数据类
/** 带单员的一笔成交,放进 leader.fills 的就是它 */
public record LeaderFill(String fillId, // 成交编号,比如 LF-1001
String leaderId, // 带单员 ID,比如 L1
String symbol, // 交易对,比如 BTC-USDT
String side, // BUY / SELL
BigDecimal qty, // 带单员这笔的数量
BigDecimal price, // 成交价
BigDecimal leaderEquity) {} // 带单员本金,按比例跟单要用
public enum Mode { RATIO, FIXED } // 按比例 / 固定金额
/** 一个跟单者的设置 */
public record Follower(String id, BigDecimal equity, Mode mode,
BigDecimal value, // RATIO 时是倍数,FIXED 时是金额 (U)
BigDecimal available) {} // 额度:最多能下多少 U
/** 要给某个跟单者下的一张单,放进 copy.orders 的就是它 */
public record CopyOrder(String clientOid, String followerId, String symbol, String side, BigDecimal qty) {}
// 下面几个小工具自己补上,都很短:
// toJson(obj) / parseFill(json) / parseOrder(json):用 Jackson 的 ObjectMapper 转换
// followerRepo.findByLeader("L1"):练习里直接返回写死的 3 个 Follower
// skip(f, 原因):打印一行"跳过 F3:太小"
static final BigDecimal LOT = new BigDecimal("0.001"); // lotSize
static final BigDecimal MIN_QTY = new BigDecimal("0.001"); // 最小下单量
第 4 步:生产者,发出"带单员成交了"
Properties p = new Properties();
p.put("bootstrap.servers", "localhost:9092");
p.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("acks", "all"); // 所有同步中的副本都写成功才算发送成功(单机只有 1 个副本)
p.put("enable.idempotence", "true"); // 生产者自己重试时,不会往 Kafka 里写出重复消息
try (KafkaProducer<String, String> producer = new KafkaProducer<>(p)) {
LeaderFill fill = new LeaderFill("LF-1001", "L1", "BTC-USDT", "BUY",
new BigDecimal("1"), new BigDecimal("60000"), new BigDecimal("10000"));
// key = 带单员 ID:同一个带单员的成交进同一个分区,按顺序处理
producer.send(new ProducerRecord<>("leader.fills", fill.leaderId(), toJson(fill))).get(); // .get():等 Kafka 确认收到
}
| 配置项 | 大白话 |
|---|---|
bootstrap.servers | Kafka 在哪:地址和端口,本机就是 localhost:9092 |
key.serializer / value.serializer | key 和消息内容怎么变成字节发出去。这里两者都是字符串,消息内容是 JSON 字符串 |
acks | Kafka 收到后要写好几份副本才回"成功"。all = 所有同步中的副本都写好,最稳,稍慢 |
enable.idempotence | 生产者网络抖动自己重发时,Kafka 不会存两份。注意:它只管"发"这一侧,第 8 节说的"消费者重复处理"它管不了,那个要靠 clientOid |
group.id | 消费者组的名字(第 4 节)。同名的消费者一起分担分区,书签也按组记 |
enable.auto.commit | 要不要每隔几秒自动夹书签。关掉,改成处理完手动 commitSync(),才能控制"先处理完,再夹书签" |
auto.offset.reset | 这个组还没有书签时从哪读。earliest = 从最早的消息,latest = 只读启动之后新来的 |
第 5 步:第一段,扇出服务
Properties c = new Properties();
c.put("bootstrap.servers", "localhost:9092");
c.put("group.id", "copy-fanout"); // 消费者组的名字:同名的消费者一起分担分区
c.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
c.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
c.put("enable.auto.commit", "false"); // 关掉自动提交:全部处理完才手动夹书签
c.put("auto.offset.reset", "earliest"); // 这个组第一次启动、还没有书签时,从最早的消息读起
// p 就是第 4 步那份生产者配置:扇出服务既是消费者(读 leader.fills),又是生产者(写 copy.orders)
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(c);
KafkaProducer<String, String> producer = new KafkaProducer<>(p)) {
consumer.subscribe(List.of("leader.fills"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(200)); // 最多等 200 毫秒,拿一批消息
for (ConsumerRecord<String, String> r : records) {
LeaderFill fill = parseFill(r.value());
for (Follower f : followerRepo.findByLeader(fill.leaderId())) {
BigDecimal qty = CopySizer.size(fill, f, LOT); // 第 7 节的规则
if (qty.compareTo(MIN_QTY) < 0) { skip(f, "太小"); continue; }
if (qty.multiply(fill.price()).compareTo(f.available()) > 0) { skip(f, "额度不够"); continue; }
String clientOid = fill.fillId() + "-" + f.id(); // 幂等编号:重试多少次都一样
// key = 跟单者 ID:第二段按跟单者分区,流量均匀
producer.send(new ProducerRecord<>("copy.orders", f.id(),
toJson(new CopyOrder(clientOid, f.id(), fill.symbol(), fill.side(), qty))));
}
}
producer.flush(); // 确认这一批全部发出去了
consumer.commitSync(); // 再夹书签。顺序反过来就可能丢单
}
}
第 6 步:第二段,下单消费者(用一个 Set 假装交易所)
c.put("group.id", "copy-order-worker"); // 其余配置和第 5 步一样,只是换一个组名
Set<String> seen = new HashSet<>(); // 假装交易所:记住收到过的 clientOid
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(c)) {
consumer.subscribe(List.of("copy.orders"));
while (true) {
for (ConsumerRecord<String, String> r : consumer.poll(Duration.ofMillis(200))) {
CopyOrder o = parseOrder(r.value());
if (!seen.add(o.clientOid())) { // 同一个编号第二次来:拒绝
System.out.println("重复订单被拒绝: " + o.clientOid());
continue;
}
// 真实系统在这里调交易所的下单接口;练习里打印出来就够了
System.out.println("下单: " + o.clientOid() + " qty=" + o.qty());
}
consumer.commitSync();
}
}
第 7 步:验证
followerRepo 里给 L1 放 3 个跟单者,额度都给够(别让第 7 节的规则把人跳过)。先启动第二段和第一段,再运行生产者发 1 条,第二段应该打印 3 行"下单"。producer.flush() 之后、commitSync() 之前加一行 if (!records.isEmpty()) System.exit(1);(只在真的处理了消息时才退出)。② 第二段保持运行,别重启它,它的
seen 集合在内存里,一重启就忘了。③ 带着这一行重启扇出服务,发一条新成交(比如 LF-1002)。扇出服务处理完、发出 3 张单,还没提交书签就退出了。
④ 删掉那一行,再启动扇出服务。它会从上一个书签开始,把 LF-1002 再处理一遍,第二段应该打印 3 行"重复订单被拒绝"。
UUID.randomUUID().toString(),用一条新成交(LF-1003)重复 ①–④。这次第二段会打印 6 行"下单",同一笔成交给每个跟单者下了两次单。词典和自测
今天用到的词都在这里,包括前几天学过、今天又用到的。
自测 7 题
回到 七日冲刺 D6:先看两段式设计图(现在每个框你都能讲清楚了),再在分区动画里比较三种分区键、触发一次 rebalance,最后做 copySizes 编程题。然后按第 10 节在本机跑一遍。