跳到主要内容

消息骨架:Queue 抽象与持久化

30 秒导读: Kestra 的 Executor、Worker、Scheduler 三个角色互不直接调用,它们之间的每一次 "我把活儿交给你"都变成一条消息,丢进一条队列。这一章讲的就是这条底层消息总线:一套 QueueInterface 接口定义"消息怎么进、怎么出",一个基于数据库的实现把消息存进一张 queues, 靠轮询把消息发给订阅者。关键设计:接口只有一套,实现只有 JDBC 一种,单机跑就把数据源换成 内存里的 H2,集群跑就换成 Postgres/MySQL——换库不换代码

本章只讲传输与存储(消息怎么排队、怎么落库、怎么发出去)。消息内容被谁怎么处理,是 执行引擎WorkerScheduler 各章的事。


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 + queueNamecore/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 / subscribeBroadcastqueue-jdbc/client/JdbcQueueClient.java:33
JdbcQueueFactory为每种消息类型造一个队列 Beanqueue-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一条消息只被一个消费者拿到,拿走即删ExecutionWorkerTaskResult
Broadcast 广播BroadcastQueueInterface每个消费者都收到全部消息ExecutionKilledFlowInterface(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):

参数默认含义
minPollInterval25ms高峰期最快轮询间隔
maxPollInterval500ms空闲期最慢轮询间隔
pollSwitchInterval60s空轮多久后切到最慢档
pollSize100一次最多捞多少条
immediateRepolltrue拉到货就立刻再轮一次

一句话:队列空闲时省数据库,繁忙时抢低延迟。 这是"轮询"这种朴素机制能撑生产的关键调优。

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,当 typeh2memory 时,给出一个内存 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)数据源适用重启后消息
memoryjdbc:h2:mem:public本地开发、测试丢失
h2H2(可文件)单机小规模视配置
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. 巧妙之处(可带走的技术)

  1. 一表多队列 + type 列。 不为每条队列建表,而是一张 queues 表用 type 区分。运维简单 (一套迁移、一套索引),代价是所有队列共享一张热表。(baseline-queue-h2.sql)

  2. Dispatch 用 FOR UPDATE SKIP LOCKED 把数据库当消息中间件。 不引 Kafka/RabbitMQ,靠关系库 原生的跳锁能力实现并发竞争消费,让"零外部依赖"的单机部署成为可能。 (JdbcQueueClient.java:188)

  3. 广播用"每订阅者游标 + 定时清道夫"而非删行。 读写分离干净:消费者只推自己的游标,回收交给 独立的 @Scheduled 任务,避免消费者互相踩删。(JdbcQueueClient.java:278JdbcQueueCleaner.java:51)

  4. 自适应轮询退避。 用"maxInterval 反复除 2 逼近 minInterval"的阶梯,在数据库成本和消息延迟 间自动平衡;拉满一批就立刻重轮抢低延迟。(QueuePoller.java:36QueuePollerConfiguration.java:36)

  5. 内存与集群共用一份实现。 "memory 模式 = 内存 H2 数据源",开发体验与生产行为同构,是极省心 智负担的设计取舍。(DatasourceProvider.java:19)

  6. 失败即回滚即重投。 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.javaQueueInterface
新接口根core/src/main/java/io/kestra/core/queues/GenericQueueInterface.javaGenericQueueInterfaceaddListenerqueueName
四种投递语义接口core/src/main/java/io/kestra/core/queues/DispatchQueueInterface.javaDispatchQueueInterfaceBroadcastQueueInterfaceKeyedDispatchQueueInterfaceVNodeDispatchQueueInterface
消息标记接口core/src/main/java/io/kestra/core/queues/event/Event.javaEvent.key
key 计算core/src/main/java/io/kestra/core/queues/QueueService.javaQueueService.key
序列化 + 消息保护 + vnodequeue/src/main/java/io/kestra/queue/QueueService.javaserializedeserializecomputeVNode
队列抽象基类queue/src/main/java/io/kestra/queue/AbstractQueue.javaAbstractQueuequeueNametrackSubscriberclose
Dispatch 抽象queue/src/main/java/io/kestra/queue/AbstractDispatchQueue.javaemitdoEmitsubscriber
Broadcast/Keyed/VNode 抽象queue/src/main/java/io/kestra/queue/Abstract{Broadcast,KeyedDispatch,VNodeDispatch}Queue.javaemitdoEmit
订阅者基类(pause/close/致命错误)queue/src/main/java/io/kestra/queue/AbstractSubscriber.javaprocessMessagewaitIfPausedmarkEndclose
轮询订阅循环queue/src/main/java/io/kestra/queue/AbstractPollingSubscriber.javainternalSubscribepollpollBatch
退避轮询queue/src/main/java/io/kestra/queue/poller/QueuePoller.javapollOnce
轮询配置(阶梯计算)queue/src/main/java/io/kestra/queue/poller/QueuePollerConfiguration.javacomputeStepsStep
队列 Bean 注解queue/src/main/java/io/kestra/queue/QueueBean.java@QueueBean
队列全局配置(前缀/消息保护)queue/src/main/java/io/kestra/queue/QueueConfiguration.javaQueueConfigurationMessageProtection
JDBC 队列工厂(消息↔语义绑定)queue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcQueueFactory.javaJdbcQueueFactory
JDBC Dispatch/Broadcast/Keyed/VNode 队列queue-jdbc/src/main/java/io/kestra/queue/jdbc/Jdbc*Queue.javaJdbcDispatchQueueJdbcBroadcastQueueJdbcKeyedDispatchQueueJdbcVNodeDispatchQueue
SQL 核心(publish/竞争取/广播取)queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/JdbcQueueClient.javapublishsubscribeDispatchsubscribeBroadcastfetchMaxOffsetqueueLag
Dispatch/Broadcast 订阅者queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/Jdbc{Dispatch,Broadcast}Subscriber.javapollinitmaxOffset
广播清道夫queue-jdbc/src/main/java/io/kestra/queue/jdbc/client/JdbcQueueCleaner.javadeleteQueue
依赖聚合 recordqueue-jdbc/src/main/java/io/kestra/queue/jdbc/JdbcDependencies.javaJdbcDependencies
队列启用条件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.javaJdbcQueueConfiguration
队列行记录jdbc/src/main/java/io/kestra/jdbc/JdbcQueueItem.javaJdbcQueueItem
内存模式数据源(= H2 mem)repository-memory/src/main/java/io/kestra/runner/memory/DatasourceProvider.javaDatasourceProviderH2RepositoryOrQueue
队列表 DDLjdbc-h2/src/main/resources/migrations/baseline-queue-h2.sqlqueues
仓储 Bean 注解core/src/main/java/io/kestra/core/repositories/RepositoryBean.java@RepositoryBean
仓储 JDBC 基类jdbc/src/main/java/io/kestra/jdbc/repository/AbstractJdbcRepository.javaAbstractJdbcRepositoryfielddefaultFilter
仓储末端子类示例jdbc-h2/src/main/java/io/kestra/repository/h2/H2FlowRepository.javaH2FlowRepository

相邻章节: 领域模型 · 执行引擎 · Worker 与 RunContext · Scheduler 与 Trigger · 插件系统与 Web 层