feat: schedule due processing from next dueAt instead of polling

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
0264408
2026-09-20 16:54:13 +08:00
co-authored by Cursor
parent 8d08a2840b
commit 02557a9b86
15 changed files with 183 additions and 12 deletions
+2 -1
View File
@@ -111,7 +111,8 @@ ordo:
definitions: definitions:
location: classpath*:ordo/*.json # 默认值;启动时 publish 加载 location: classpath*:ordo/*.json # 默认值;启动时 publish 加载
due: due:
poll-ms: 0 # >0 时轮询 processDue poll-ms: 0 # >0 启用调度;值为最长空闲,按下次 dueAt 唤醒,满批续拉
batch-size: 100
rest: rest:
enabled: false enabled: false
base-path: /ordo base-path: /ordo
+5 -2
View File
@@ -75,7 +75,8 @@ ordo:
definitions: definitions:
location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 publish location: classpath*:ordo/*.json # 启动时对每个 JSON 调用 publish
due: due:
poll-ms: 0 # >0 时轮询 processDue;默认不调度 poll-ms: 0 # >0 启用调度;值为最长空闲,按下次 dueAt 唤醒,满批续拉
batch-size: 100
rest: rest:
enabled: false enabled: false
base-path: /ordo base-path: /ordo
@@ -262,7 +263,8 @@ ordo.cancel(instance.id(), "admin", "政策变更");
- `approve` / `reject` / `reassign`:`actor` 必须等于该任务当前 `assignee`,否则 `UnauthorizedTaskOperationException`。 - `approve` / `reject` / `reassign`:`actor` 必须等于该任务当前 `assignee`,否则 `UnauthorizedTaskOperationException`。
- `reassign`:仅 `PENDING` 任务;同一任务 id,办理人改为 `newAssignee`,不推进步骤。`newAssignee` 不可空白、不可等于当前 `assignee`,且同一步不能已有该人的 `PENDING` 任务,否则 `IllegalArgumentException`。不经过 `AssigneeResolver`。人工转派不改 `dueAt`。 - `reassign`:仅 `PENDING` 任务;同一任务 id,办理人改为 `newAssignee`,不推进步骤。`newAssignee` 不可空白、不可等于当前 `assignee`,且同一步不能已有该人的 `PENDING` 任务,否则 `IllegalArgumentException`。不经过 `AssigneeResolver`。人工转派不改 `dueAt`。
- `processDue(limit)`:认领 `dueAt <= now` 的 PENDING 任务(`limit > 0`),按步上 `due.then` 执行:`reassign` 换办理人(`to` 走 `AssigneeResolver`)、`notify` 可选 `ActionHandler`、`goto` 跳过当前步 PENDING 并进入 `to` 步骤(`to` 必须是步骤 id)。每种策略对一张任务最多成功一次(清空 `dueAt`)。引擎无后台线程;Spring 下 `ordo.due.poll-ms > 0` 才轮询。 - `processDue(limit)`:认领 `dueAt <= now` 的 PENDING 任务(`limit > 0`),按步上 `due.then` 执行:`reassign` 换办理人(`to` 走 `AssigneeResolver`)、`notify` 可选 `ActionHandler`、`goto` 跳过当前步 PENDING 并进入 `to` 步骤(`to` 必须是步骤 id)。每种策略对一张任务最多成功一次(清空 `dueAt`)。引擎无后台线程;Spring 下 `ordo.due.poll-ms > 0` 才调度:按下次 `dueAt` 唤醒,满批续拉,`poll-ms` 为最长空闲。
- `nextDueAt()`:PENDING 且仍有 `dueAt` 的最早到期时刻;没有则 empty。
- 任务非 `PENDING`:`TaskAlreadyCompletedException`。 - 任务非 `PENDING`:`TaskAlreadyCompletedException`。
- `withdraw`:仅 `initiator`,否则 `UnauthorizedInstanceOperationException`;实例非 `RUNNING`:`InstanceAlreadyCompletedException`。 - `withdraw`:仅 `initiator`,否则 `UnauthorizedInstanceOperationException`;实例非 `RUNNING`:`InstanceAlreadyCompletedException`。
- `cancel`:`actor` 非空即可,**不校验**是否发起人;实例须为 `RUNNING`,否则 `InstanceAlreadyCompletedException`。谁能调用由宿主决定。 - `cancel`:`actor` 非空即可,**不校验**是否发起人;实例须为 `RUNNING`,否则 `InstanceAlreadyCompletedException`。谁能调用由宿主决定。
@@ -392,6 +394,7 @@ ACTION 成功事件发生在提交之后,因此排在同轮事务内写入的
| `approve` / `reject` | 办理当前 PENDING 任务 | | `approve` / `reject` | 办理当前 PENDING 任务 |
| `reassign` | 当前办理人把 PENDING 任务转给他人 | | `reassign` | 当前办理人把 PENDING 任务转给他人 |
| `processDue` | 认领并处理已到期 PENDING 任务 | | `processDue` | 认领并处理已到期 PENDING 任务 |
| `nextDueAt` | 下一笔 PENDING 任务的 `dueAt` |
| `withdraw` | 发起人撤回 | | `withdraw` | 发起人撤回 |
| `cancel` | 管理员/系统取消(引擎不鉴权角色) | | `cancel` | 管理员/系统取消(引擎不鉴权角色) |
| `find*` | 按 id / 待办索引读取 | | `find*` | 按 id / 待办索引读取 |
@@ -11,6 +11,7 @@ import com.jetlumen.ordo.api.runtime.ProcessContext;
import com.jetlumen.ordo.api.runtime.ProcessEvent; import com.jetlumen.ordo.api.runtime.ProcessEvent;
import com.jetlumen.ordo.api.runtime.ProcessInstance; import com.jetlumen.ordo.api.runtime.ProcessInstance;
import java.time.Instant;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
@@ -40,6 +41,9 @@ public interface OrdoEngine {
/** Claims and processes up to {@code limit} overdue pending tasks. {@code limit} must be positive. */ /** Claims and processes up to {@code limit} overdue pending tasks. {@code limit} must be positive. */
int processDue(int limit); int processDue(int limit);
/** Earliest {@code dueAt} among pending tasks, or empty if none are scheduled. */
Optional<Instant> nextDueAt();
default ProcessInstance withdraw(String instanceId, String actor) { default ProcessInstance withdraw(String instanceId, String actor) {
return withdraw(instanceId, actor, null); return withdraw(instanceId, actor, null);
} }
@@ -43,6 +43,9 @@ public interface ApprovalTaskRepository {
List<ApprovalTask> findDuePending(Instant now, int limit); List<ApprovalTask> findDuePending(Instant now, int limit);
/** Earliest {@code dueAt} among pending tasks, or empty if none are scheduled. */
Optional<Instant> findNextDueAt();
boolean claimIfDue(String taskId, String expectedAssignee, Instant now); boolean claimIfDue(String taskId, String expectedAssignee, Instant now);
/** Paginated, filterable query; results are ordered newest-first (created_at desc). */ /** Paginated, filterable query; results are ordered newest-first (created_at desc). */
@@ -345,6 +345,11 @@ public final class DefaultOrdoEngine implements OrdoEngine {
return processed; return processed;
} }
@Override
public synchronized Optional<Instant> nextDueAt() {
return taskRepository.findNextDueAt();
}
private boolean escalateDueTask(ApprovalTask overdue, Instant now, List<PendingAction> queued, private boolean escalateDueTask(ApprovalTask overdue, Instant now, List<PendingAction> queued,
List<ProcessEvent> events) { List<ProcessEvent> events) {
if (!taskRepository.claimIfDue(overdue.id(), overdue.assignee(), now)) { if (!taskRepository.claimIfDue(overdue.id(), overdue.assignee(), now)) {
@@ -23,6 +23,7 @@ import com.jetlumen.ordo.core.repository.InMemoryProcessHistoryRepository;
import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository; import com.jetlumen.ordo.core.repository.InMemoryProcessInstanceRepository;
import java.time.Clock; import java.time.Clock;
import java.time.Instant;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
@@ -181,4 +182,9 @@ public final class InMemoryOrdoEngine implements OrdoEngine {
public int processDue(int limit) { public int processDue(int limit) {
return delegate.processDue(limit); return delegate.processDue(limit);
} }
@Override
public Optional<Instant> nextDueAt() {
return delegate.nextDueAt();
}
} }
@@ -113,6 +113,15 @@ public final class InMemoryApprovalTaskRepository implements ApprovalTaskReposit
.toList(); .toList();
} }
@Override
public synchronized Optional<Instant> findNextDueAt() {
return tasks.values().stream()
.filter(task -> task.status() == TaskStatus.PENDING)
.map(ApprovalTask::dueAt)
.filter(Objects::nonNull)
.min(Comparator.naturalOrder());
}
@Override @Override
public synchronized boolean claimIfDue(String taskId, String expectedAssignee, Instant now) { public synchronized boolean claimIfDue(String taskId, String expectedAssignee, Instant now) {
ApprovalTask current = tasks.get(taskId); ApprovalTask current = tasks.get(taskId);
@@ -5,24 +5,35 @@ import org.springframework.context.SmartLifecycle;
import java.lang.System.Logger; import java.lang.System.Logger;
import java.lang.System.Logger.Level; import java.lang.System.Logger.Level;
import java.time.Clock;
import java.time.Instant;
import java.util.Objects; import java.util.Objects;
import java.util.concurrent.Executors; import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
/** Optional poller that calls {@link OrdoEngine#processDue(int)} when {@code ordo.due.poll-ms} is positive. */ /** Optional scheduler that drains {@link OrdoEngine#processDue(int)} when {@code ordo.due.poll-ms} is positive. */
public final class OrdoDuePoller implements SmartLifecycle { public final class OrdoDuePoller implements SmartLifecycle {
private static final Logger LOG = System.getLogger("ordo"); private static final Logger LOG = System.getLogger("ordo");
private static final int BATCH_SIZE = 100;
private final OrdoEngine ordoEngine; private final OrdoEngine ordoEngine;
private final Clock clock;
private final long pollMs; private final long pollMs;
private final int batchSize;
private ScheduledExecutorService executor; private ScheduledExecutorService executor;
private ScheduledFuture<?> scheduled;
private volatile boolean running; private volatile boolean running;
private volatile boolean ticking;
public OrdoDuePoller(OrdoEngine ordoEngine, long pollMs) { public OrdoDuePoller(OrdoEngine ordoEngine, Clock clock, long pollMs, int batchSize) {
this.ordoEngine = Objects.requireNonNull(ordoEngine, "ordoEngine must not be null"); this.ordoEngine = Objects.requireNonNull(ordoEngine, "ordoEngine must not be null");
this.clock = Objects.requireNonNull(clock, "clock must not be null");
this.pollMs = pollMs; this.pollMs = pollMs;
if (batchSize <= 0) {
throw new IllegalArgumentException("batch size must be positive");
}
this.batchSize = batchSize;
} }
@Override @Override
@@ -40,18 +51,51 @@ public final class OrdoDuePoller implements SmartLifecycle {
thread.setDaemon(true); thread.setDaemon(true);
return thread; return thread;
}); });
executor.scheduleWithFixedDelay(this::tick, pollMs, pollMs, TimeUnit.MILLISECONDS);
running = true; running = true;
schedule(0);
}
void wake() {
if (!running || ticking) {
return;
}
schedule(0);
} }
private void tick() { private void tick() {
ticking = true;
try { try {
ordoEngine.processDue(BATCH_SIZE); int processed;
do {
processed = ordoEngine.processDue(batchSize);
} while (processed == batchSize);
} catch (RuntimeException e) { } catch (RuntimeException e) {
LOG.log(Level.WARNING, "processDue failed", e); LOG.log(Level.WARNING, "processDue failed", e);
} finally {
if (running) {
schedule(nextDelayMs());
}
ticking = false;
} }
} }
private long nextDelayMs() {
Instant now = clock.instant();
return ordoEngine.nextDueAt()
.map(dueAt -> Math.min(Math.max(0L, dueAt.toEpochMilli() - now.toEpochMilli()), pollMs))
.orElse(pollMs);
}
private synchronized void schedule(long delayMs) {
if (!running || executor == null) {
return;
}
if (scheduled != null) {
scheduled.cancel(false);
}
scheduled = executor.schedule(this::tick, delayMs, TimeUnit.MILLISECONDS);
}
@Override @Override
public void stop() { public void stop() {
running = false; running = false;
@@ -59,6 +103,7 @@ public final class OrdoDuePoller implements SmartLifecycle {
executor.shutdownNow(); executor.shutdownNow();
executor = null; executor = null;
} }
scheduled = null;
} }
@Override @Override
@@ -0,0 +1,25 @@
package com.jetlumen.ordo.spring;
import com.jetlumen.ordo.api.runtime.ProcessEvent;
import com.jetlumen.ordo.api.runtime.ProcessEventType;
import com.jetlumen.ordo.api.spi.OrdoEventListener;
/** Forwards {@code TASK_CREATED} to {@link OrdoDuePoller} after the engine listener snapshot is taken. */
public final class OrdoDueWakeBridge implements OrdoEventListener {
private volatile OrdoDuePoller poller;
public void attach(OrdoDuePoller poller) {
this.poller = poller;
}
@Override
public void onEvent(ProcessEvent event) {
if (event.type() != ProcessEventType.TASK_CREATED) {
return;
}
OrdoDuePoller current = poller;
if (current != null) {
current.wake();
}
}
}
@@ -51,9 +51,12 @@ public class OrdoProperties {
} }
public static class Due { public static class Due {
/** Poll interval in milliseconds. {@code 0} disables scheduling. */ /** Enables scheduling when positive; also the maximum idle sleep in milliseconds. */
private long pollMs; private long pollMs;
/** Tasks claimed per {@code processDue} call while draining. */
private int batchSize = 100;
public long getPollMs() { public long getPollMs() {
return pollMs; return pollMs;
} }
@@ -61,6 +64,14 @@ public class OrdoProperties {
public void setPollMs(long pollMs) { public void setPollMs(long pollMs) {
this.pollMs = pollMs; this.pollMs = pollMs;
} }
public int getBatchSize() {
return batchSize;
}
public void setBatchSize(int batchSize) {
this.batchSize = batchSize;
}
} }
public static class Jdbc { public static class Jdbc {
@@ -19,6 +19,7 @@ import com.jetlumen.ordo.core.spi.DispatchingActionHandler;
import com.jetlumen.ordo.core.spi.DispatchingRoutingCondition; import com.jetlumen.ordo.core.spi.DispatchingRoutingCondition;
import com.jetlumen.ordo.spring.OrdoDefinitionLoader; import com.jetlumen.ordo.spring.OrdoDefinitionLoader;
import com.jetlumen.ordo.spring.OrdoDuePoller; import com.jetlumen.ordo.spring.OrdoDuePoller;
import com.jetlumen.ordo.spring.OrdoDueWakeBridge;
import com.jetlumen.ordo.spring.OrdoProperties; import com.jetlumen.ordo.spring.OrdoProperties;
import com.jetlumen.ordo.storage.jdbc.JdbcActionExecutionRepository; import com.jetlumen.ordo.storage.jdbc.JdbcActionExecutionRepository;
import com.jetlumen.ordo.storage.jdbc.JdbcApprovalTaskRepository; import com.jetlumen.ordo.storage.jdbc.JdbcApprovalTaskRepository;
@@ -152,6 +153,12 @@ public class OrdoJdbcAutoConfiguration {
return new JdbcInstanceTokenRepository(connectionProvider, ordoSqlDialect); return new JdbcInstanceTokenRepository(connectionProvider, ordoSqlDialect);
} }
@Bean
@ConditionalOnMissingBean
public OrdoDueWakeBridge ordoDueWakeBridge() {
return new OrdoDueWakeBridge();
}
@Bean @Bean
@ConditionalOnMissingBean @ConditionalOnMissingBean
public OrdoEngine ordoEngine(Clock ordoClock, public OrdoEngine ordoEngine(Clock ordoClock,
@@ -181,7 +188,11 @@ public class OrdoJdbcAutoConfiguration {
@Bean @Bean
@ConditionalOnMissingBean @ConditionalOnMissingBean
public OrdoDuePoller ordoDuePoller(OrdoEngine ordoEngine, OrdoProperties ordoProperties) { public OrdoDuePoller ordoDuePoller(OrdoEngine ordoEngine, Clock ordoClock, OrdoProperties ordoProperties,
return new OrdoDuePoller(ordoEngine, ordoProperties.getDue().getPollMs()); OrdoDueWakeBridge ordoDueWakeBridge) {
OrdoDuePoller poller = new OrdoDuePoller(ordoEngine, ordoClock, ordoProperties.getDue().getPollMs(),
ordoProperties.getDue().getBatchSize());
ordoDueWakeBridge.attach(poller);
return poller;
} }
} }
@@ -11,6 +11,7 @@ import com.jetlumen.ordo.storage.jdbc.dialect.SqlDialects;
import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper; import com.jetlumen.ordo.storage.jdbc.mapper.ApprovalTaskMapper;
import java.sql.Connection; import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.PreparedStatement; import java.sql.PreparedStatement;
import java.sql.ResultSet; import java.sql.ResultSet;
import java.sql.SQLException; import java.sql.SQLException;
@@ -18,6 +19,7 @@ import java.sql.Timestamp;
import java.time.Instant; import java.time.Instant;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Locale;
import java.util.Objects; import java.util.Objects;
import java.util.Optional; import java.util.Optional;
@@ -52,6 +54,8 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
private static final String SELECT_DUE_PENDING_BASE = private static final String SELECT_DUE_PENDING_BASE =
"SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL" "SELECT " + TASK_COLUMNS + " FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL"
+ " AND due_at <= ? ORDER BY due_at, id"; + " AND due_at <= ? ORDER BY due_at, id";
private static final String SELECT_NEXT_DUE_AT =
"SELECT MIN(due_at) FROM ordo_approval_task WHERE status = 'PENDING' AND due_at IS NOT NULL";
private static final String TASK_COLUMNS_QUALIFIED = private static final String TASK_COLUMNS_QUALIFIED =
"t.id, t.instance_id, t.step_id, t.task_name, t.assignee, t.status, t.created_at, t.completed_at," "t.id, t.instance_id, t.step_id, t.task_name, t.assignee, t.status, t.created_at, t.completed_at,"
+ " t.action_actor, t.action_comment, t.action_at, t.due_at"; + " t.action_actor, t.action_comment, t.action_at, t.due_at";
@@ -67,7 +71,7 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) { public JdbcApprovalTaskRepository(JdbcConnectionProvider connectionProvider, SqlDialect dialect) {
this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null"); this.connectionProvider = Objects.requireNonNull(connectionProvider, "connectionProvider must not be null");
this.dialect = Objects.requireNonNull(dialect, "dialect must not be null"); this.dialect = Objects.requireNonNull(dialect, "dialect must not be null");
this.selectDuePending = dialect.limit(SELECT_DUE_PENDING_BASE, false); this.selectDuePending = duePendingSql(this.connectionProvider, this.dialect);
} }
@Override @Override
@@ -180,6 +184,23 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
} }
} }
@Override
public Optional<Instant> findNextDueAt() {
Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(SELECT_NEXT_DUE_AT);
ResultSet resultSet = select.executeQuery()) {
if (!resultSet.next()) {
return Optional.empty();
}
Timestamp timestamp = resultSet.getTimestamp(1);
return timestamp == null ? Optional.empty() : Optional.of(timestamp.toInstant());
} catch (SQLException e) {
throw new JdbcStorageException("failed to query next due at", e);
} finally {
connectionProvider.close(connection);
}
}
@Override @Override
public boolean claimIfDue(String taskId, String expectedAssignee, Instant now) { public boolean claimIfDue(String taskId, String expectedAssignee, Instant now) {
Objects.requireNonNull(taskId, "taskId must not be null"); Objects.requireNonNull(taskId, "taskId must not be null");
@@ -275,6 +296,24 @@ public final class JdbcApprovalTaskRepository implements ApprovalTaskRepository
} }
} }
private static String duePendingSql(JdbcConnectionProvider connectionProvider, SqlDialect dialect) {
String limited = dialect.limit(SELECT_DUE_PENDING_BASE, false);
return supportsSkipLocked(connectionProvider) ? dialect.forUpdateSkipLocked(limited) : limited;
}
private static boolean supportsSkipLocked(JdbcConnectionProvider connectionProvider) {
Connection connection = connectionProvider.getConnection();
try {
DatabaseMetaData metaData = connection.getMetaData();
String product = metaData.getDatabaseProductName();
return product != null && !product.toLowerCase(Locale.ROOT).contains("h2");
} catch (SQLException e) {
return false;
} finally {
connectionProvider.close(connection);
}
}
private List<ApprovalTask> findAll(String sql, String parameter) { private List<ApprovalTask> findAll(String sql, String parameter) {
Connection connection = connectionProvider.getConnection(); Connection connection = connectionProvider.getConnection();
try (PreparedStatement select = connection.prepareStatement(sql)) { try (PreparedStatement select = connection.prepareStatement(sql)) {
@@ -16,4 +16,9 @@ abstract class LimitOffsetSqlDialect implements SqlDialect {
public final String forUpdate(String sql) { public final String forUpdate(String sql) {
return sql + " FOR UPDATE"; return sql + " FOR UPDATE";
} }
@Override
public final String forUpdateSkipLocked(String sql) {
return sql + " FOR UPDATE SKIP LOCKED";
}
} }
@@ -25,5 +25,8 @@ public interface SqlDialect {
String forUpdate(String sql); String forUpdate(String sql);
/** Appends {@code FOR UPDATE SKIP LOCKED}. */
String forUpdateSkipLocked(String sql);
String[] flywayLocations(); String[] flywayLocations();
} }
@@ -54,6 +54,7 @@ class SqlDialectsTest {
assertEquals("SELECT 1 LIMIT ? OFFSET ?", dialect.limit("SELECT 1")); assertEquals("SELECT 1 LIMIT ? OFFSET ?", dialect.limit("SELECT 1"));
assertEquals("SELECT 1 LIMIT ?", dialect.limit("SELECT 1", false)); assertEquals("SELECT 1 LIMIT ?", dialect.limit("SELECT 1", false));
assertEquals("SELECT 1 FOR UPDATE", dialect.forUpdate("SELECT 1")); assertEquals("SELECT 1 FOR UPDATE", dialect.forUpdate("SELECT 1"));
assertEquals("SELECT 1 FOR UPDATE SKIP LOCKED", dialect.forUpdateSkipLocked("SELECT 1"));
assertEquals("classpath:db/postgresql/migration", dialect.flywayLocations()[0]); assertEquals("classpath:db/postgresql/migration", dialect.flywayLocations()[0]);
} }