消息骨架:Queue 抽象与持久化
30 秒导读: Kestra 的 Executor、Worker、Scheduler 三个角色互不直接调用,它们之间的每一次 "我把活儿交给你"都变成一条消息,丢进一条队列。这一章讲的就是这条底层消息总线:一套
QueueInterface接口定义"消息怎么进、怎么出",一个基于数据库的实现把消息存进一张queues表, 靠轮询把消息发给订阅者。关键设计:接口只有一套,实现只有 JDBC 一种,单机跑就把数据源换成 内存里的 H2,集群跑就换成 Postgres/MySQL——换库不换代码。
本章只讲传输与存储(消息怎么排队、怎么落库、怎么发出去)。消息内容被谁怎么处理,是 执行引擎、Worker、 Scheduler 各章的事。
1. 这是什么(零基础也能懂)
一句话定义: 队列(Queue)是一根异步消息管道——一端有人往里塞消息(emit / 发布),另一端 有人从里取消息(receive / 订阅),两端在时间和进程上都可以解耦。
它解决什么问题? 设想 Kestra 在跑一个工作流:
- Executor 算出"下一步该运行任务 A"——它不该直接
new Worker().run(A),因为 Worker 可能在 另一台机器上,可能正忙,可能刚崩了。 - 正确做法:Executor 把"任务 A 待运行"写成一条消息丢进队列,转身去处理别的执行;某个空闲的 Worker 从队列捞出这条消息去干,干完再把结果丢回另一条队列,Executor 再捞回来推进状态。
于是队列成了三个角色之间唯一的耦合点。角色之间不认识彼此,只认识队列:
写消息 读消息
Executor ──emit──▶ ┌─────────┐ ──receive──▶ Worker
Worker ──emit──▶ │ 队列 │ ──receive──▶ Executor
Scheduler──emit──▶ └─────────┘ ──receive──▶ ...
(生产者) 总线 (消费者)
一句话直觉: 把队列当成办公室里的收件筐。谁有活儿就写张纸条丢 进筐(emit),谁有空就从筐里 拿一张来做(receive)。筐满不满、谁来拿,发纸条的人不用管。Kestra 的特别之处是:这个"收件筐" 不是内存里的一个链表,而是数据库里的一张表——这样即使进程重启、换机器,筐里的纸条还在。
用起来什么样? 对业务代码,队列就是一个注入进来的 Bean,两个动作:
// 示意,非源码:生产者一端
executionQueue.emit(execution); // 把一条执行丢进队列
// 消费者一端:注册一个回调,消息来了就调它
QueueSubscriber<Execution> sub = executionQueue.subscriber();
sub.subscribe(either -> process(either.getLeft())); // 每来一条消息回调一次
2. 顶层全景(它大概怎么转)
2.1 三层结构
从上到下三层:接口层(在 core 模块,定义契约)→ 抽象基类层(在 queue 模块,把通用逻辑
写死)→ JDBC 实现层(在 queue-jdbc 模块,真正读写数据库)。
怎么读这张图:越往下越具体。上层定"做什么",下层定"怎么用数据库做"。
┌──────────────────────────────────────────────────────────┐
│ 接口层 core/queues/*Interface.java │
│ GenericQueueInterface ── addListener / queueName │
│ ├─ DispatchQueueInterface 竞争消费(每条只发一次)│
│ ├─ BroadcastQueueInterface 广播(人人 都收到) │
│ ├─ KeyedDispatchQueueInterface 带路由键的竞争消费 │
│ └─ VNodeDispatchQueueInterface 按虚拟节点分片 │
└───────────────────────────┬──────────────────────────────┘
│ implements
┌───────────────────────────▼──────────────────────────────┐
│ 抽象层 queue/Abstract*.java(通用:序列化/计数/生命周期) │
│ AbstractQueue │
│ ├─ AbstractDispatchQueue ├─ AbstractBroadcastQueue│
│ ├─ AbstractKeyedDispatchQueue └─ AbstractVNodeDispatch │
│ AbstractSubscriber ← AbstractPollingSubscriber(轮询循环)│
└───────────────────────────┬──────────────────────────────┘
│ extends
┌───────────────────────────▼──────────────────────────────┐
│ 实现层 queue-jdbc/*.java(唯一实现) │
│ JdbcDispatchQueue / JdbcBroadcastQueue / JdbcKeyed... │
│ └── JdbcQueueClient ──── 读写 ────▶ 表 queues │
└───────────────────────────┬──────────────────────────────┘
│ 数据源可换
┌──────────┬──────────┼──────────┬──────────┐
H2内存 H2文件 Postgres MySQL (同一套 SQL)
memory h2
2.2 部件一句话职责
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
QueueInterface | 队列的老接口:emit/receive/delete 的全套签名 | core/queues/QueueInterface.java:11 |
GenericQueueInterface | 新接口族的根:只留 addListener + queueName | core/queues/GenericQueueInterface.java:7 |
Dispatch/Broadcast/Keyed/VNode QueueInterface | 四种投递语义的接口 | core/queues/*QueueInterface.java |
QueueService(core) | 给消息对象算 key(去重/路由用) | core/queues/QueueService.java:10 |
QueueService(queue) | 序列化/反序列化 + 消息大小保护 + vnode 计算 | queue/QueueService.java:65 |
AbstractQueue | 通用:队列名、监听器、指标、订阅者跟踪、关闭 | queue/AbstractQueue.java:24 |
AbstractPollingSubscriber | 通用轮询循环:反复调 poll() 直到关闭 | queue/AbstractPollingSubscriber.java:26 |
QueuePoller | 单次轮询 + 退避(backoff)计算 | queue/poller/QueuePoller.java:14 |
JdbcQueueClient | 真正的 SQL:publish / subscribeDispatch / subscribeBroadcast | queue-jdbc/client/JdbcQueueClient.java:33 |
JdbcQueueFactory | 为每种消息类型造一个队列 Bean | queue-jdbc/JdbcQueueFactory.java:32 |
JdbcQueueCleaner | 定时清理广播队列的过期消息 | queue-jdbc/client/JdbcQueueCleaner.java:29 |
表 queues | 所有队列消息的落地表(一表多队列) | jdbc-h2/.../baseline-queue-h2.sql |
2.3 主线走一遍(高层)
一条消息的完整旅程:
emit(msg)
│ ① QueueService.serialize:把对象 → JSON 字节,顺便查大小是否超限
│ ② JdbcQueueClient.publish:INSERT 一行进 queues 表(type=队列名, value=JSON)
▼
[ 消息静静躺在 queues 表里,带一个自增 offset ]
▲
│ ③ 订阅者的轮询循环每隔几十毫秒执行一次 poll()
│ ④ JdbcQueueClient.subscribeXxx:SELECT 一批出来
│ ⑤ QueueService.deserialize:JSON 字节 → 对象
│ ⑥ 回调业务 consumer
receive(consumer)
3. 核心原理(逐个机制,由浅入深)
3.1 一套接口,四种投递语义
同一根总线,但"消息该发给谁"有四种不同规矩。这是理解整章的第一把钥匙。所有消息类型都实现
一个标记接口 Event,它只要求一个方法 key()(消息的业务主键,用于日志/去重):
// core/queues/event/Event.java
public interface Event {
String key();
}
四种语义:
| 语义 | 接口 | 规矩 | 典型消息 |
|---|---|---|---|
| Dispatch 竞争消费 | DispatchQueueInterface | 一条消息只被一个消费者拿到,拿走即删 | Execution、WorkerTaskResult |
| Broadcast 广播 | BroadcastQueueInterface | 每个消费者都收到全部消息 | ExecutionKilled、FlowInterface(flow 变更) |
| KeyedDispatch 带键竞争 | KeyedDispatchQueueInterface | 竞争消费,但按 routingKey 分流 | WorkerJobEvent(按 worker group 路由) |
| VNodeDispatch 虚拟节点分片 | VNodeDispatchQueueInterface | 按 key 哈希到虚拟节点,竞争消费 | TriggerEvent(Scheduler 分片) |
为什么要分 Dispatch 和 Broadcast?因为语义天差地别:一次任务只能有一个 Worker 去做(Dispatch, 否则重复执行);而"某个流被删了"这种事件必须让所有 Scheduler/Executor 都知道(Broadcast)。 两种语义在数据库层的实现完全不同——见 §3.4。
JdbcQueueFactory 就是把每一种业务消息绑定到某一种语义。例如:
// queue-jdbc/JdbcQueueFactory.java:35 executionQueue —— Execution 用 Dispatch
public DispatchQueueInterface<Execution> executionQueue(JdbcDependencies d) {
return new JdbcDispatchQueue<>(Execution.class, ...);
}
// :61 killQueue —— ExecutionKilled 用 Broadcast
public BroadcastQueueInterface<ExecutionKilled> killQueue(JdbcDependencies d) {
return new JdbcBroadcastQueue<>(ExecutionKilled.class, ...);
}
3.2 emit 路径:序列化 + 消息保护
生产者调 emit,抽象基类先做三件通用事(以 Dispatch 为例):
// queue/AbstractDispatchQueue.java:23 emit
public final void emit(T message) throws QueueException {
this.doEmit(this.queueService.serialize(this.cls, message), message.key()); // ① 序列化 ② 交给子类落库
listeners().forEach(l -> l.accept(message)); // ③ 同步通知本地监听器(测试用)
this.emitCounter.increment(); // 指标 +1
}
doEmit 是抽象方法,由 JdbcDispatchQueue 实现成一次 INSERT:
// queue-jdbc/JdbcDispatchQueue.java:48 doEmit
protected void doEmit(byte[] message, String key) throws QueueException {
jdbcQueueClient.publish(this.queueName(), null, key, new String(message));
}
精华在序列化这一步的"消息保护"。 QueueService.serialize 把对象转成 JSON 字节后,会检查大小:
// queue/QueueService.java:65 serialize(节选)
byte[] serialize = MAPPER.writeValueAsBytes(message);
if (messageProtection.enabled && serialize.length >= messageProtection.limit) {
// 但对"已终止的 Execution"网开一面,让它过去
if (!(message instanceof Execution e) || !e.getState().isTerminated()) {
throw new MessageTooBigException(...); // 太大就拒发
}
}
为什么?数据库的一行放不下无限大的消息;一条超大消息会拖垮整条队列。所以 Kestra 设了默认阈值,
超了就抛 MessageTooBigException,提示你把大输出存进内部存储而不是塞进消息。唯一的例外是
"已经结束的执行"必须让它落库,否则执行状态就永远推进不下去了。
3.3 订阅 = 一个不停轮询的循环
Kestra 的队列没有"消息推送(push)"机制——它靠轮询(poll):订阅者起一个后台线程,反复问
数据库"有新消息吗?"。这套循环写在 AbstractPollingSubscriber:
// queue/AbstractPollingSubscriber.java:60 internalSubscribe(节选)
while (this.isActive()) {
this.waitIfPaused(); // 若被暂停,在这里阻塞
lastPoll = queuePoller.pollOnce(lastPoll, steps); // 执行一次轮询查询
}
为什么不会把数据库问爆? 因为 QueuePoller.pollOnce 有一套自适应退避:有消息就快轮,
没消息就慢慢拉长间隔。
// queue/poller/QueuePoller.java:36 pollOnce(节选)
Integer count = pollingQuery.call(); // 查一次,返回捞到几条
if (count > 0) {
sleep = configuration.minPollInterval(); // 有货 → 用最短间隔
if (count.equals(configuration.pollSize())) // 一次拉满 → 立刻再轮,不睡
return lastPoll;
} else {
// 没货 → 根据"空轮了多久"逐级切换到更长的间隔
sleep = ...selectedSteps.getLast().pollInterval();
}
Thread.sleep(sleep);
间隔的档位由 QueuePollerConfiguration.computeSteps 预先算好:从 maxPollInterval 出发,不断
除以 2 逼近 minPollInterval,形成一个由慢到快的阶梯。默认值(JdbcQueueConfiguration):
| 参数 | 默认 | 含义 |
|---|---|---|
minPollInterval | 25ms | 高峰期最快轮询间隔 |
maxPollInterval | 500ms | 空闲期最慢轮询间隔 |
pollSwitchInterval | 60s | 空轮多久后切到最慢档 |
pollSize | 100 | 一次最多捞多少条 |
immediateRepoll | true | 拉到货就立刻再轮一次 |
一句话:队列空闲时省数据库,繁忙时抢低延迟。 这是"轮询"这种朴素机制能撑生产的关键调优。
3.4 Dispatch vs Broadcast 落库取舍
这是整章工程含量最高的一支。两种语义共用一张 queues 表、共用轮询框架,但"从表里取消息"
的 SQL 逻辑截然相反。
Dispatch(竞争消费):抢锁 + 删行。 多个消费者可能同时来抢同一批消息,必须保证一条消息只被
一个人拿走。用的是 FOR UPDATE SKIP LOCKED(选出行并加锁,别人绕开已锁的行):
// queue-jdbc/client/JdbcQueueClient.java:170 subscribeDispatch(节选,单事务内)
var result = context.select(OFFSET, VALUE).from(table)
.where(TYPE.eq(queue))
.orderBy(OFFSET.asc())
.limit(pollSize)
.forUpdate().skipLocked() // ← 关键:抢到就锁住,别的消费者跳过
.fetch();
// 逐条回调 consumer,然后:
context.delete(table).where(OFFSET.in(processedItems)).execute(); // ← 处理完删掉
要点:①SKIP LOCKED 让多个消费者并发无冲突地各拿各的一批;②消息处理完即删除,表里只留
"还没被消费"的消息;③整批在一个事务里,回调抛异常则事务回滚 → 消息不删 → 下次重投(at-least-once)。
Broadcast(广播):游标 + 不删。 每个订阅者要看到全部消息,所以不能删——删了别人就看不到了。
改用每个订阅者自己记一个 maxOffset 游标,每次只查"offset 比我记的更大"的行:
// queue-jdbc/client/JdbcQueueClient.java:268 subscribeBroadcast(节选)
var select = context.select(OFFSET, VALUE).from(table).where(TYPE.eq(queue));
if (maxOffset != null) select = select.and(OFFSET.gt(maxOffset)); // ← 只看比游标新的
var result = select.orderBy(OFFSET.asc()).limit(pollSize).fetch(); // 注意:没有 FOR UPDATE,没有 delete
// 回调后把游标推到本批最大 offset
maxOffsetResult = result.stream().map(r -> r.get(OFFSET)).max(...).orElse(null);
订阅者启动时从 fetchMaxOffset 把游标初始化到"当前最大 offset",意味着广播订阅者只收开始订阅
之后的新消息(不回放历史):
// queue-jdbc/client/JdbcBroadcastSubscriber.java:46 init
protected void init() { maxOffset = jdbcQueueClient.fetchMaxOffset(queueName); markReady(); }
那广播消息谁来删?没有消费者删,于是 有个独立的定时清道夫 JdbcQueueCleaner,按保留期(默认 1
小时)扫掉过期的广播消息:
// queue-jdbc/client/JdbcQueueCleaner.java:51 deleteQueue
@Scheduled(fixedDelay = "${kestra.jdbc.queue.cleaner.fixed-delay:1h}")
public long deleteQueue() { /* 对每个广播队列 DELETE created <= now - retention */ }
两种语义对照:
| 维度 | Dispatch(竞争消费) | Broadcast(广播) |
|---|---|---|
| 谁收到 | 只有一个消费者 | 每个消费者都收到 |
| 并发保证 | FOR UPDATE SKIP LOCKED | 无锁,靠各自游标 |
| 消息删除 | 消费成功后立即删 | 消费者不删,靠 JdbcQueueCleaner 按保留期清 |
| 是否回放历史 | 不涉及(拿走就没了) | 否,只收订阅后的新消息(游标初始化到 max) |
| 位置记录 | 无(表即状态) | 每个订阅者内存里存 maxOffset |
| 失败重投 | 事务回滚 → 不删 → 重投 | 游标未推进 → 下次重读 |
KeyedDispatch 与 VNode 是 Dispatch 的变体。 它们复用同一套 subscribeDispatch,只是给消息带上
routing_key,查询时按键过滤。VNode 的键是对消息 key 做哈希算出的虚拟节点号,用于把
TriggerEvent 均匀分片到多个 Scheduler:
// queue-jdbc/JdbcVNodeDispatchQueue.java:33 doEmit
jdbcQueueClient.publish(queueName(),
this.vNodeRoutingKey(this.queueService.computeVNode(key)), // key → vnode_N
key, new String(message));
3.5 订阅者的生命周期:pause / resume / close
AbstractSubscriber 用一组原子标志 + 一把锁管理订阅者状态,支撑"优雅暂停/关闭":
- pause/resume: 暂停不是停线程,而是让轮询循环在
waitIfPaused()处用Condition.await()阻塞;resume()用signalAll()唤醒(AbstractSubscriber.java:153、:213)。用途:优雅关闭、 维护窗口、背压。 - close: 用
CountDownLatch等轮询线程真正退出,超时 30 秒(AbstractSubscriber.java:260)。 - 致命错误 → 关整个进程: 若消费循环抛出不可恢复异常(非事务/连接类),调
markEnd(cause)记录并触发KestraContext.shutdown()——故意不 ack 消息,好让它重启后重投给别的实例 (AbstractSubscriber.java:248)。这是"宁可重启也不吞消息"的取舍。
一个巧妙细节:轮询循环会吞掉关库时的事务/连接异常,因为那是进程正常关闭时数据源被关的正常现象,
不该当致命错误(AbstractPollingSubscriber.java:82)。
4. 两种落地:内存版 vs JDBC 版,其实是同一套
任务里常把 Kestra 说成"内存队列 vs JDBC 队列两种实现"。但读源码后要诚实纠正:在本 commit,队列 只有一种实现——JDBC。 所谓"内存模式",是把 JDBC 的数据源换成内存里的 H2,而非另写一套内存队列。
证据一:标记注解把四种模式都指向同一实现。 激活 JDBC 队列的条件是
kestra.queue.type ∈ {memory, h2, mysql, postgres},四个值走同一份代码:
// queue-jdbc/JdbcQueueEnabled.java
@Requires(property = "kestra.queue.type", pattern = "memory|h2|mysql|postgres")
public @interface JdbcQueueEnabled {}
证据二:memory 就是内存 H2。 repository-memory 模块(内部包名 runner.memory)只有一个类
DatasourceProvider,当 type 是 h2 或 memory 时,给出一个内存 H2 数据源:
// repository-memory/.../DatasourceProvider.java:19
memory.setUrl("jdbc:h2:mem:public"); // ← "内存队列" = H2 in-memory 数据库
证据三:全库只有 JDBC 系的具体队列类。 grep "extends Abstract*Queue" 的非测试实现,只有
queue-jdbc 下的 JdbcDispatchQueue / JdbcBroadcastQueue / JdbcKeyedDispatchQueue / JdbcVNodeDispatchQueue——没有任何 MemoryQueue。
所以真正的主设计线是:
一套 QueueInterface + 一套 Jdbc 实现 + 一张 queues 表
│
kestra.queue.type 只换「数据源」这一件事
┌──────────────┬──────────────┬──────────────┬──────────────┐
memory / h2mem h2(文件) postgres mysql
单机、重启即丢 单机、可持久 集群、生产 集群、生产
模式(kestra.queue.type) | 数据源 | 适用 | 重启后消息 |
|---|---|---|---|
memory | jdbc:h2:mem:public | 本地开发、测试 | 丢失 |
h2 | H2(可文件) | 单机小规模 | 视配置 |
postgres / mysql | 独立数据库 | 集群、生产 | 保留 |
repository-memory 依赖 jdbc + jdbc-h2(见其 build.gradle),runner-memory 又只依赖
repository-memory——内存模块 = H2 模块的一层薄封装,坐实了"没有独立内存实现"这一点。
这条设计的价值:开发者
docker run一个零依赖的 Kestra(内存 H2)体验到的行为,和生产集群 (Postgres)完全一致,因为跑的是同一份队列代码、同一套 SQL。测试与生产不再是两套心智模型。
5. 持久化:queues 表结构与"瞬态"语义
所有队列的消息都落在同一张表 queues,靠 type 列区分是哪条队列。表结构
(jdbc-h2/.../baseline-queue-h2.sql,Postgres/MySQL 同构):
CREATE TABLE queues (
"offset" BIGINT NOT NULL AUTO_INCREMENT PRIMARY KEY, -- 全局自增序号,兼作广播游标
"type" VARCHAR(250) NOT NULL, -- 队列名(= 消息类名转下划线),一表多队列靠它区分
"routing_key" VARCHAR(250), -- Keyed/VNode 的路由键,普通队列为 NULL
"key" VARCHAR(250) NOT NULL, -- 消息业务主键(Event.key())
"value" TEXT NOT NULL, -- 消息体 JSON(Postgres 里是 JSONB)
"created" TIMESTAMP NOT NULL -- 落库时间,清道夫按它算保留期
);
注:
type列在2.0.09迁移里从INT改成VARCHAR(250)(2.0.09-queue-index-and-type-h2.sql),与JdbcQueueClient.TYPE = field("type", String.class)一致。行的形状对应记录JdbcQueueItem,但读路径只SELECT offset, value,消息类型信息在type列而非反序列化时用。
索引专为两种访问模式建:
queues_type__key__offset (type, routing_key, offset)—— Dispatch 按队列+路由键顺序取。queues_created__type (created, type)—— 清道夫按时间清广播消息。
关键语义:队列是"瞬态"的。 迁移脚本里写得很直白:
"The queue is transient: in-flight messages are lost on restart and replayed from executions state." (
2.0.05-queue-h2.sql注释)
也就是说,queues 表存的是在途消息,不是权威状态。权威状态是 executions 等业务表
(见 第 1 章 领域模型)。万一在途消息丢了,系统靠执行状态重新推导该发什么
消息。这解释了为什么内存 H2 模式"重启即丢"是可接受的——真相不在队列里。
发布时的一个"数据脏值"防护:Postgres 的 JSONB 拒绝某些非法 Unicode(如 �、孤立代理项),
JdbcQueueClient 会把这类 DataException 转成可恢复的 UnsupportedMessageException,让
Worker 能优雅处理而不是让整条队列卡死(JdbcQueueClient.java:102)。
6. 仓储层模式:与队列同构的"三层可插拔"
队列负责在途消息,仓储(Repository)负责权威状态的读写。两者用的是同一套可插拔套路, 理解了队列就理解了仓储。
三层结构(和 §2.1 队列如出一辙):
core/repositories/XxxRepositoryInterface.java ← 接口(业务只依赖它)
▲ implements
jdbc/repository/AbstractJdbcXxxRepository.java ← 通用 JDBC 实现(jOOQ 写 SQL)
▲ extends
jdbc-h2/…/H2XxxRepository.java ← 每种数据库一个薄子类
jdbc-postgres/…/PostgresXxxRepository.java (只处理方言差异)
jdbc-mysql/…/MysqlXxxRepository.java
以 Flow 仓储为例,末端子类只是绑定数据源和方言:
// jdbc-h2/.../H2FlowRepository.java:23
@RepositoryBean
@H2RepositoryEnabled
public class H2FlowRepository extends AbstractJdbcFlowRepository { ... }
两个约定 Bean 注解把"何时装配"声明化:
@RepositoryBean=@Singleton+ "只在非 Worker 服务上激活" (core/repositories/RepositoryBean.java)。呼应 Worker 约束:Worker 不直接碰仓储,只经 MetaStore/StateStore 门面。@QueueBean=@Singleton+@Bean(preDestroy="close"),保证队列关闭时优雅释放订阅者 (queue/QueueBean.java)。
AbstractJdbcRepository 提供所有仓储共享的基建:field(name) 造 jOOQ 字段、defaultFilter()
统一加上"未删除 + 租户隔离"的过滤条件(jdbc/repository/AbstractJdbcRepository.java:40)。
一句话:队列和仓储是同一枚硬币的两面——同一套"core 定接口、jdbc 写通用实现、每种库出一个薄 子类"的插拔模式,一边管消息流动,一边管状态存储。换数据库时,两边一起从 H2 切到 Postgres,业务 代码零改动。
7. 巧妙之处(可带走的技术)
-
一表多队列 +
type列。 不为每条队列建表,而是一张queues表用type区分。运维简单 (一套迁移、一套索引),代价是所有队列共享一张热表。(baseline-queue-h2.sql) -
Dispatch 用
FOR UPDATE SKIP LOCKED把数据库当消息中间件。 不引 Kafka/RabbitMQ,靠关系库 原生的跳锁能力实现并发竞争消费,让"零外部依赖"的单机部署成为可能。 (JdbcQueueClient.java:188) -
广播用"每订阅者游标 + 定时清道夫"而非删行。 读写分离干净:消费者只推自己的游标,回收交给 独立的
@Scheduled任务,避免消费者互相踩删。(JdbcQueueClient.java:278、JdbcQueueCleaner.java:51) -
自适应轮询退避。 用"maxInterval 反复除 2 逼近 minInterval"的阶梯,在数据库成本和消息延迟 间自动 平衡;拉满一批就立刻重轮抢低延迟。(
QueuePoller.java:36、QueuePollerConfiguration.java:36) -
内存与集群共用一份实现。 "memory 模式 = 内存 H2 数据源",开发体验与生产行为同构,是极省心 智负担的设计取舍。(
DatasourceProvider.java:19) -
失败即回滚即重投。 Dispatch 的取/删在一个事务里,消费回调抛异常则整批不删,天然 at-least-once;致命错误宁可重启进程也不 ack,把消息留给别的实例。(
AbstractSubscriber.java:248)
8. 边界与局限(诚实)
- 只有 JDBC 一种实现。 本 commit 没有独立的内存/Kafka 队列;所有模式都落到关系库。吞吐上限
受单张
queues热表 + 轮询频率约束,不是为超高吞吐(百万 msg/s)场景设计的。 - 推送靠轮询,有固有延迟。 最快 25ms 一轮(默认),空闲期可拉到 500ms;不是事件驱动的即时推送。
- 广播只收订阅之后的消息。 游标初始化到当前 max offset,订阅者启动前的历史广播消息看不到
(
JdbcBroadcastSubscriber.java:46)。 - 消息有大小上限。 超过
message-protection.limit会被拒(终止态执行除外);大 payload 需自行 存进内部存储(QueueService.java:78)。 - 队列是瞬态的,不是审计日志。 消费即删(Dispatch)或定时清理(Broadcast);要追溯得看
executions等业务表。 memory模式重启丢消息。 生产必须用 Postgres/MySQL;内存 H2 仅供开发测试。
9. 代码地图(导航索引)
| 主题 | 文件路径 | 关键符号 |
|---|---|---|
| 老队列接口(emit/receive 全套) | core/src/main/java/io/kestra/core/queues/QueueInterface.java | QueueInterface |
| 新接口根 | core/src/main/java/io/kestra/core/queues/GenericQueueInterface.java | GenericQueueInterface、addListener、queueName |
| 四种投递语义接口 | core/src/main/java/io/kestra/core/queues/DispatchQueueInterface.java 等 | DispatchQueueInterface、BroadcastQueueInterface、KeyedDispatchQueueInterface、VNodeDispatchQueueInterface |
| 消息标记接口 | core/src/main/java/io/kestra/core/queues/event/Event.java | Event.key |
| key 计算 | core/src/main/java/io/kestra/core/queues/QueueService.java | QueueService.key |
| 序列化 + 消息保护 + vnode | queue/src/main/java/io/kestra/queue/QueueService.java | serialize、deserialize、computeVNode |
| 队列抽象基类 | queue/src/main/java/io/kestra/queue/AbstractQueue.java | AbstractQueue、queueName、trackSubscriber、close |
| Dispatch 抽象 | queue/src/main/java/io/kestra/queue/AbstractDispatchQueue.java | emit、doEmit、subscriber |
| Broadcast/Keyed/VNode 抽象 | queue/src/main/java/io/kestra/queue/Abstract{Broadcast,KeyedDispatch,VNodeDispatch}Queue.java | emit、doEmit |
| 订阅者基类(pause/close/致命错误) | queue/src/main/java/io/kestra/queue/AbstractSubscriber.java | processMessage、waitIfPaused、markEnd、close |
| 轮询订阅循环 | queue/src/main/java/io/kestra/queue/AbstractPollingSubscriber.java | internalSubscribe、poll、pollBatch |
| 退避轮询 | queue/src/main/java/io/kestra/queue/poller/QueuePoller.java | pollOnce |
| 轮询配置(阶梯计算) | queue/src/main/java/io/kestra/queue/poller/QueuePollerConfiguration.java | computeSteps、Step |
| 队列 Bean 注解 | queue/src/main/java/io/kestra/queue/QueueBean.java | @QueueBean |
| 队列全局配置(前缀/消息保护) | queue/src/main/java/io/kestra/queue/QueueConfiguration.java | QueueConfiguration、MessageProtection |
| JDBC 队列工厂(消息↔语义绑定) | queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcQueueFactory.java | JdbcQueueFactory |
| JDBC Dispatch/Broadcast/Keyed/VNode 队列 | queue-jdbc/src/main/java/io/kestra/queue/jdbc/Jdbc*Queue.java | JdbcDispatchQueue、JdbcBroadcastQueue、JdbcKeyedDispatchQueue、JdbcVNodeDispatchQueue |
| SQL 核心(publish/竞争取/广播取) | queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/JdbcQueueClient.java | publish、subscribeDispatch、subscribeBroadcast、fetchMaxOffset、queueLag |
| Dispatch/Broadcast 订阅者 | queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/Jdbc{Dispatch,Broadcast}Subscriber.java | poll、init、maxOffset |
| 广播清道夫 | queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/JdbcQueueCleaner.java | deleteQueue |
| 依赖聚合 record | queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcDependencies.java | JdbcDependencies |
| 队列启用条件 | queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcQueueEnabled.java | @Requires(kestra.queue.type = memory|h2|mysql|postgres) |
| 轮询默认值 | jdbc/src/main/java/io/kestra/jdbc/runner/JdbcQueueConfiguration.java | JdbcQueueConfiguration |
| 队列行记录 | jdbc/src/main/java/io/kestra/jdbc/JdbcQueueItem.java | JdbcQueueItem |
| 内存模式数据源(= H2 mem) | repository-memory/src/main/java/io/kestra/runner/memory/DatasourceProvider.java | DatasourceProvider、H2RepositoryOrQueue |
| 队列表 DDL | jdbc-h2/src/main/resources/migrations/baseline-queue-h2.sql | 表 queues |
| 仓储 Bean 注解 | core/src/main/java/io/kestra/core/repositories/RepositoryBean.java | @RepositoryBean |
| 仓储 JDBC 基类 | jdbc/src/main/java/io/kestra/jdbc/repository/AbstractJdbcRepository.java | AbstractJdbcRepository、field、defaultFilter |
| 仓储末端子类示例 | jdbc-h2/src/main/java/io/kestra/repository/h2/H2FlowRepository.java | H2FlowRepository |
相邻章节: 领域模型 · 执行引擎 · Worker 与 RunContext · Scheduler 与 Trigger · 插件系统与 Web 层