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

D6 零基础速成笔记

把 D1 讲过的"跟单"真正用 Kafka 跑起来。Kafka 也从零讲:它是什么、消息怎么排队、为什么会重复。页面里带下划虚线的词,点一下就能看解释。

怎么读按顺序读,一共 11 节,大约 2–3 小时。第 10 节是动手部分,可以边读边在本机敲代码。
每个概念都分四层大白话 → 生活例子 → 算一遍 → Java 类比。前两层看懂就够了。
读完以后回到七日冲刺的 D6,看两段式设计图、分区键动画,做 copySizes 编程题。
1

今天要做什么:跟单从头到尾

复习的词跟单 · 带单员 · 跟单者 · 分润

今天把 D1 第 11 节讲的跟单真正用代码跑起来。先把业务流程完整走一遍,后面每一节都会对应到其中某一步。

一步一步看一笔跟单带单员 L1,3 个跟单者

    录屏:滚到这里会自动一步步播放,每步停 4 秒左右,高亮的是当前这一步。想细看就点"暂停",用"上一步 / 下一步"来回翻。

    第 2–6 节讲 Kafka 怎么把"成交 → 找人 → 下单"这几步连起来;第 7 节讲"每人下多少";第 8 节讲下单重试时最容易踩的坑。

    2

    为什么要用消息队列

    本节新词消息 · 消息队列 · 生产者 · 消费者

    消息队列是一个"先放着、慢慢处理"的中转站。发的一方把消息放进去就走,处理的一方按自己的速度一条条取出来处理。

    消息队列 message queue

    生活例子
    餐厅点单:服务员把点菜单插到后厨的单子架上,马上去招呼下一桌;厨师按顺序一张张做。服务员不用站在灶台边等菜做好。单子架就是消息队列,每张点菜单就是一条消息。
    算一遍
    带单员成交 1 笔,要给 3000 个跟单者下单,每单 5 毫秒,全部做完要 15 秒。如果"成交服务"直接调用"下单服务"等它做完,就要卡 15 秒。改成把一条消息放进队列,成交服务 1 毫秒就返回,下单的活交给专门的程序慢慢干。
    Java 类比
    和 BlockingQueue 一样:一边 put(),一边 take()。区别是 Kafka 是独立部署的服务,消息存在磁盘上,程序重启了消息还在,多台机器都能来读。

    生产者 和 消费者 producer / consumer

    大白话
    生产者:往队列里放消息的程序,这里是发出"带单员成交了"的那个服务。
    消费者:从队列里取消息来处理的程序,这里是跟单服务。
    生活例子
    服务员是生产者,厨师是消费者。
    3

    Kafka 的 topic、分区和 key

    本节新词Kafka · topic · 分区 · key · 同分区有序

    Kafka 是最常用的消息队列之一,你们组用得很多。先记住三个词:topic 是"哪一类消息",分区是"这类消息分成几条队来排",key 决定"这条消息排哪一队"。

    topic 主题

    大白话
    一类消息的名字,可以理解成一个"频道"。比如 放"带单员成交了"(leader = 带单员,fills = 成交), 放"要给跟单者下的单"(copy = 跟单,orders = 订单)。名字是我们自己起的,Kafka 不规定。
    生活例子
    银行取号机上的"个人业务"和"对公业务",是两种不同的号。
    Java 类比
    一个有名字的队列,Map<String, Queue> 里的那个 String。

    分区 partition

    大白话
    一个 topic 拆成几条并行的队伍,每条队伍同一时间只由一个消费者处理。拆开是为了让多个消费者同时干活。
    生活例子
    "个人业务"开了 4 个窗口,每个窗口前排一条队。
    算一遍
    每秒来 1200 条消息,一个消费者每秒能处理 400 条,那至少要 3 个分区、3 个消费者才处理得过来。
    Java 类比
    List<BlockingQueue> partitions,每个队列配一个线程。

    key 和 同分区有序

    大白话
    发消息时带一个 key。Kafka 用 key 算一个哈希值,再对分区数取余,决定它进哪一队。同一个 key 永远进同一队;同一队里先来先处理,所以同一个 key 的消息一定按顺序处理。不同队之间谁先谁后不保证。
    生活例子
    银行规定"同一家公司的人都去 2 号窗口",这家公司的业务就一定按到达顺序办。
    算一遍
    4 个分区,key 是 "L1",算出来余数是 2 → 进 2 号分区。L1 的每一笔成交都在 2 号分区排队,开仓一定比平仓先处理。
    Java 类比
    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 里按顺序排着。

    4

    消费者组、offset 和 lag

    本节新词消费者组 · offset · 提交 · lag

    消费者组 consumer group

    大白话
    一组一起干活的消费者。Kafka 把分区分给组里的成员,每个分区同一时间只给一个成员。
    生活例子
    4 个窗口,组里 4 个柜员:一人一个窗口。只有 2 个柜员:每人管 2 个窗口。来了 6 个柜员:有 2 个没窗口可坐,只能闲着。
    算一遍
    4 个分区、6 个消费者 → 只有 4 个在干活。分区数决定了最多几个消费者能同时干活。
    Java 类比
    线程池里的线程,但每个线程绑定了固定的几个队列。

    offset 和 提交 offset / commit

    大白话
    offset 是每条消息在分区里的序号,从 0 开始。消费者处理完后告诉 Kafka"这个分区我处理到第几条了",这叫提交 offset。
    生活例子
    看书夹书签。下次打开,从书签那页接着看。
    算一遍
    分区里有第 0–99 条。消费者处理完第 59 条,提交 60(意思是下次从 60 开始)。这时它挂了,接手的消费者从 60 开始读,不漏也不重复。
    但如果它处理到第 79 条才挂,书签还停在 60,那 60–79 这 20 条会被再处理一遍。第 8 节专门讲这个。
    Java 类比
    consumer.commitSync() 就是夹书签。
    常见误区 · 打开自动提交图省事Kafka 可以设成每隔几秒自动夹书签(第 10 节配置表里的 enable.auto.commit)。它只看时间,不看你处理完没有:书签可能已经挪到第 80 条,你才处理到第 70 条,这时程序挂了,第 70–79 条就被跳过了,也就是丢消息。练习里一律关掉自动提交,处理完再手动 commitSync():宁可重复(第 8 节有防线),不能丢。

    lag 消费延迟、堆积量

    大白话
    堆着还没处理的消息有多少条。= 分区里一共有多少条 − 书签(已提交的 offset)的位置。
    生活例子
    窗口前还排着几个人。
    接着上面的例子
    分区里有第 0–99 条,一共 100 条;书签在 60。lag = 100 − 60 = 40,也就是第 60–99 这 40 条还没处理。
    算一遍
    每秒来 1200 条,每秒只能处理 1000 条,每秒多堆 200 条。1 分钟后 lag = 12000 条,新来的一条要等 12000 ÷ 1000 = 12 秒才轮到。对跟单来说,这 12 秒价格可能已经走出去很远了。
    Java 类比
    队列的 size()。
    lag 计算器跑 60 秒会堆多少
    真正在干活的消费者–
    每秒总处理能力–
    60 秒后 lag–
    新消息要排队–

    怎么玩:默认只有 2 个消费者,每秒处理 800 条、进来 1200 条,60 秒后堆了 24,000 条。把"消费者数"拉到 3,堆积消失;再拉到 6,看"真正在干活的消费者"停在 4,因为分区只有 4 个,多出来的人闲着。

    5

    分区键怎么选

    本节新词分区键 · 热点分区

    分区键就是发消息时拿哪个字段当 key。选哪个字段,决定了两件事:谁和谁的消息一定按顺序处理,以及流量会不会都挤进同一个分区。

    用什么当 key保证什么顺序问题
    带单员 ID同一个带单员的成交按顺序处理(开仓一定在平仓之前)热门带单员流量大,他所在的分区会堆积
    交易对同一个交易对的消息按顺序跟单用不上这个顺序;而且 BTC 流量最大,同样会挤
    跟单者 ID同一个跟单者的订单按顺序流量均匀。但前提是消息已经拆成了"每个跟单者一条"(第 6 节)

    热点分区 hot partition

    大白话
    某一个分区的消息特别多,它的消费者忙不过来;其他分区却很闲。
    生活例子
    一家大公司的人都被规定去 2 号窗口,2 号排长队,其他窗口没人。
    算一遍
    4 个分区,每个消费者每秒处理 400 条,总流量每秒 1200 条,平均每个分区 300 条,本来够用。但热门带单员 L1 一个人就占了 600 条,全进 P2,P2 每秒至少多堆 200 条;其他分区分剩下的 600 条,每个只有 150 条左右,连能力的一半都用不到。
    Java 类比
    HashMap 里大量 key 撞到同一个桶,那个桶的链表特别长。
    常见误区 · 热点分区靠加消费者解决一个分区同一时间只给一个消费者(第 4 节)。P2 堵了,再加 10 个消费者也只是闲着,P2 还是那一个人在干。办法只有两个:换分区键,或者像第 6 节那样把活拆开,摊到多个分区上。
    换一个分区键看看–

    怎么玩:在下拉框里依次切换三种 key,每种看几秒。条越长 = 这个分区堆的越多;变红 = 堆积超过 40 条。每个分区的消费者每一拍处理 4 条,总共进来 12 条,平均分的话完全够用。按带单员或交易对时,会有一条一直变长变红(热点分区);按跟单者时,四条都很短。

    6

    两段式扇出

    本节新词扇出 · 两段式

    扇出 fan-out

    大白话
    一条消息变成 N 个任务。
    生活例子
    群主发了一条通知,要一个个私聊 3000 个群员。
    算一遍
    带单员有 3000 个跟单者,每下一单 5 毫秒,一个线程逐个下完要 15 秒。排最后的人比第一个人晚 15 秒成交,这就是跟单滑点的主要来源。
    Java 类比
    for (Follower f : followers) placeOrder(f);

    两段式

    大白话
    把"拆"和"做"分成两步、放到两个 topic 里。
    第一段只负责把一条带单成交拆成 N 条跟单订单(只是算数和发消息,很快);
    第二段按跟单者 ID 分区,多个消费者并行真正去下单(慢的活摊开做)。
    生活例子
    群主把 3000 人的名单分给 8 个助手,每人私聊 375 人。
    算一遍
    3000 × 5 毫秒 ÷ 8 个分区 ≈ 1.9 秒,比 15 秒快了约 8 倍。
    Java 类比
    第一段是一个 for 循环,每次只调用 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"= 流量均匀摊开。

    一个人做 vs 分给多个人做假设:价格平均每秒挪动 0.02%
    单线程:最后一人等–
    两段式:最后一人等–
    单线程:最后一人多付–
    两段式:最后一人多付–

    怎么玩:把"第二段分区数"从 1 拖到 16,看最后一人等待的时间和多付的钱怎么下降;分区数是 1 时,两段式和单线程一样慢。"多付"按一笔 1000 U 的跟单估算:等得越久,价格挪得越远。0.02%/秒只是为了演示的假设,真实行情波动时可能大得多。

    7

    每个人下多少

    本节新词按比例跟单 · 固定金额跟单 · lotSize · 向下取整 · 额度

    第一段拆消息时,要给每个跟单者算出"下多少个"。规则只有四条,但每一条都跟钱有关。

    四条规则

    按比例
    数量 = 带单员这笔的数量 × (跟单者本金 ÷ 带单员本金) × 跟单倍数。
    例:带单员本金 10000 U 买了 1 个 BTC;你本金 5000 U、倍数 1 → 1 × 0.5 × 1 = 0.5 个。
    固定金额
    数量 = 固定金额 ÷ 价格。
    例:固定跟 200 U,价格 60000 → 200 ÷ 60000 = 0.003333… 个。
    向下取整
    交易所规定数量必须是某个最小单位的整数倍,这个最小单位叫 lotSize(比如 0.001 个)。0.003333 要向下取成 0.003,不能向上取成 0.004,否则会花掉比设定更多的钱。
    跳过
    取整后小于最小下单量 → 跳过(太小);数量 × 价格超过这个人的可用额度 → 跳过(钱不够)。跳过要记录下来并通知用户,不要偷偷改成别的数量。
    Java 类比
    金额全部用 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 取整数份,再乘回去 = 向下取整
        }
    }
    改一改跟单者,实时算结果带单员:本金 10,000 U,买 1 个 BTC,价格 60,000 · lotSize 0.001 · 最小 0.001

    额度 = 这个人现在最多能下多少 U 的单(合约里已经算上了)。默认 F3 太小被跳过、F4 额度不够被跳过。试试把 F3 的固定金额改成 60(刚好够 0.001 个),或者把 F4 的额度改成 20000,看它们变成可以下单。

    8

    重复消息从哪来,怎么防

    本节新词rebalance · 至少一次 · 幂等 · clientOid

    Kafka 保证消息不丢,但不保证只处理一次。同一条"带单员成交了"被处理两遍(这就叫重复消息),跟单者就会被下两次单,仓位翻倍。这是跟单系统里最危险的 bug 之一。

    rebalance 再平衡

    大白话
    消费者组里有人加入、退出或者挂了,Kafka 把分区重新分一遍。分的过程中,相关消费者暂停干活。
    生活例子
    一个柜员下班了,他的窗口要交给别人接着办。交接那几秒,窗口暂停服务。
    算一遍
    一次 rebalance 可能停顿几百毫秒到几秒(取决于配置和组的大小)。跟单链路里这几秒没人下单,lag 会突然冲高,跟单者的滑点也跟着变大。
    Java 类比
    consumer.subscribe(topics, new ConsumerRebalanceListener(){...}),回调里能看到自己丢了哪些分区、拿到哪些分区。

    至少一次 at-least-once

    大白话
    常见用法下,Kafka 给的保证是"消息不会丢,但可能被处理不止一次"。原因:处理完了还没来得及提交 offset 就挂了(或者发生了 rebalance),接手的人从上一个书签开始读,中间那段会再处理一遍。
    生活例子
    书看到第 80 页,书签还夹在第 60 页。睡一觉起来从第 60 页开始,60–80 页看了两遍。
    算一遍
    每处理 5 条提交一次 offset,处理完第 13 条时崩溃。上次提交是在第 10 条 → 第 11–13 条会被重新处理。每条扇出给 3 个跟单者 → 9 张订单被再发一次。

    幂等 idempotent

    大白话
    同一件事做一次和做十次,结果一样。
    生活例子
    电梯按钮:按一次和按十次,电梯都只来一次。
    Java 类比
    数据库唯一索引:同一个订单号重复 insert,第二次会被拒绝。

    clientOid 客户端订单号

    大白话
    你这边给订单起的编号,交易所会记住它,同一个编号第二次提交会被拒绝。很多交易所接口都支持这种用法,但具体规则(比如只对还没完成的订单检查重复)各家不同,以文档为准。
    跟单里的常见做法
    clientOid = 带单员成交ID + "-" + 跟单者ID,比如 = 带单员成交 LF-1001 给跟单者 F7 下的单。同一笔带单成交、同一个跟单者,不管重试多少次,算出来的编号都一样,所以不会重复下单。
    反例
    用 UUID.randomUUID() 当编号:每次重试编号都不同,交易所就当成新订单,照单全收。
    Tips · 重复是常态,不是 bugrebalance、消费者重启、上游重发,都会带来重复消息。第 10 节的 enable.idempotence 只管"生产者自己重发"这一侧,管不了消费者重复处理。所以别指望"每条只处理一次",把最后一道防线放在下单编号上:clientOid = 带单员成交 ID + 跟单者 ID,重复多少次都是同一个编号。
    模拟一次崩溃后的 rebalance每条消息扇出给 3 个跟单者
    处理完且已提交处理完但没提交,会被再处理一遍还没处理
    被重新处理的消息–
    重复发出的订单–
    交易所实际多成交–

    怎么玩:默认每 5 条提交一次、第 13 条后崩溃,黄色的第 11–13 条会被再处理一遍。把"每处理几条提交一次"拖到 1,重复就没了,但每条都要和 Kafka 来回一次,更慢;再把它拖回 5,勾上"用幂等编号",重复处理还在,但"交易所实际多成交"变成 0。这说明重复很难完全避免,真正的防线在下单编号。

    9

    公平性:谁先下单

    本节新词公平性 · 批量下单 · 轮转

    3000 个跟单者,总有人先下单、有人后下单。价格一直在变,先下单的人价格更好。让谁排第一,就是公平性问题。

    公平性

    生活例子
    抢火车票:先点提交的人先拿到座位。如果系统每次都按注册顺序处理,老用户永远排在前面,新用户永远吃亏。
    算一遍
    第 6 节的例子:两段式下,排第一的人比排最后的人早约 1.9 秒。如果价格每秒挪 0.02%,最后一个人多付约 0.04%。一笔 1000 U 的跟单就是多付 0.4 U,看起来不多,但如果每次都是同一批人吃亏,就会有投诉。

    常见做法

    批量下单
    很多交易所提供"一次提交多张订单"的接口,一次请求带几十张单,来回的时间少了,整体更快,排在后面的人等得更短。
    轮转
    每一笔带单成交,从不同的人开始下单,不让同一批人永远排最后。
    其他
    随机打乱顺序、按跟单金额分组处理等。你们组具体怎么做,进组后去看代码和设计文档,这是个很好的提问点。
     fill #1:  F1  F2  F3  F4  F5
     fill #2:  F2  F3  F4  F5  F1      ← 轮转:每次换一个人排第一
     fill #3:  F3  F4  F5  F1  F2
    10

    动手:在本机跑起来

    本节新词Docker · kafka-clients

    目标:发一条"带单员成交了",看到 3 个跟单者各自收到一张订单;中途杀掉程序再启动,看到重复的订单被拒绝。

    第 1 步:起 Kafka

    大白话
    Docker 能一行命令在本机跑起一个装好的 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 步:加依赖

    大白话
    kafka-clients 是 Kafka 官方的 Java 客户端库,里面有 KafkaProducer(生产者)和 KafkaConsumer(消费者)。JSON 解析用你熟悉的 Jackson 就行,下面的代码里省略了。
    <dependency>
      <groupId>org.apache.kafka</groupId>
      <artifactId>kafka-clients</artifactId>
      <version>3.8.0</version>
    </dependency>

    第 3 步:先定义几个数据类

    大白话
    消息里要带的东西,用 record 装起来。后面的代码都用这几个名字。
    /** 带单员的一笔成交,放进 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.serversKafka 在哪:地址和端口,本机就是 localhost:9092
    key.serializer / value.serializerkey 和消息内容怎么变成字节发出去。这里两者都是字符串,消息内容是 JSON 字符串
    acksKafka 收到后要写好几份副本才回"成功"。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 行"重复订单被拒绝"。
    对照实验
    把 clientOid 换成 UUID.randomUUID().toString(),用一条新成交(LF-1003)重复 ①–④。这次第二段会打印 6 行"下单",同一笔成交给每个跟单者下了两次单。
    11

    词典和自测

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

    自测 7 题

    读完以后

    回到 七日冲刺 D6:先看两段式设计图(现在每个框你都能讲清楚了),再在分区动画里比较三种分区键、触发一次 rebalance,最后做 copySizes 编程题。然后按第 10 节在本机跑一遍。