fix: 修复回测任务恢复逻辑和分页问题

- 修复新建任务被误判为恢复任务的问题:将 lastProcessedTradeIndex 默认值从 0 改为 null
- 修复页码计算错误:统一页码从 0 开始,确保新任务的 offset 为 0
- 修复恢复任务时跳过已处理条目的逻辑
- 添加数据库迁移文件 V30,将现有新任务的索引值改为 NULL
This commit is contained in:
WrBug
2026-01-31 21:45:02 +08:00
parent ec8cfeac77
commit 9d01c120e5
4 changed files with 54 additions and 29 deletions
@@ -134,7 +134,7 @@ data class BacktestTask(
var lastProcessedTradeTime: Long? = null,
@Column(name = "last_processed_trade_index")
var lastProcessedTradeIndex: Int = 0,
var lastProcessedTradeIndex: Int? = null,
@Column(name = "processed_trade_count")
var processedTradeCount: Int = 0
@@ -107,20 +107,21 @@ class BacktestExecutionService(
logger.info("回测时间范围: ${formatTimestamp(startTime)} - ${formatTimestamp(endTime)}, " +
"初始余额: ${task.initialBalance.toPlainString()}")
// 4. 恢复机制:如果有恢复点,计算从哪一页开始
val startPage = if (task.lastProcessedTradeIndex != null && task.lastProcessedTradeIndex >= 0) {
val lastProcessedIndex = task.lastProcessedTradeIndex
val calculatedPage = (lastProcessedIndex / size) + 1
// 4. 恢复机制:如果有恢复点,计算从哪一页开始(页码从 0 开始)
val startPage = if (task.lastProcessedTradeIndex != null) {
val lastProcessedIndex = task.lastProcessedTradeIndex!!
// 计算已处理的页码(从 0 开始)
val processedPage = lastProcessedIndex / size
// 特殊情况:如果lastProcessedTradeIndex刚好是100的倍数减1(比如99,199,299...
// 说明该页已经完全处理,应该从下一页开始
val nextPage = if (lastProcessedIndex % size == size - 1) {
calculatedPage + 1
processedPage + 1
} else {
calculatedPage
processedPage
}
logger.info("恢复任务:已处理索引=$lastProcessedIndex, 从第 $nextPage 页开始")
logger.info("恢复任务:已处理索引=$lastProcessedIndex, 计算页码=$nextPage, size=$size")
nextPage
} else {
logger.info("新任务:从第0页开始")
@@ -129,13 +130,14 @@ class BacktestExecutionService(
// 5. 分页获取和处理交易数据
var currentPage = maxOf(startPage, page)
var globalIndex = if (currentPage > 1 && task.lastProcessedTradeIndex != null) {
task.lastProcessedTradeIndex + 1
// 计算下一个要处理的全局索引(用于日志和统计)
val nextGlobalIndex = if (task.lastProcessedTradeIndex != null) {
task.lastProcessedTradeIndex!! + 1
} else {
0
}
logger.info("开始分页处理:起始页=$currentPage, 起始索引=$globalIndex")
logger.info("开始分页处理:起始页=$currentPage, 下一个要处理的索引=$nextGlobalIndex")
while (true) {
// 定期从数据库重新加载任务状态,确保能及时响应停止操作
@@ -168,9 +170,20 @@ class BacktestExecutionService(
logger.info("$currentPage 页获取到 ${pageTrades.size} 条交易")
// 处理当前页的交易
var lastProcessedIndexInPage: Int? = null
for (localIndex in pageTrades.indices) {
val leaderTrade = pageTrades[localIndex]
val index = globalIndex + localIndex
// 计算当前交易在全局数据中的索引(从 0 开始)
val index = currentPage * size + localIndex
// 如果是恢复任务,跳过已处理的条目
if (task.lastProcessedTradeIndex != null && index <= task.lastProcessedTradeIndex!!) {
logger.debug("跳过已处理的交易: index=$index, lastProcessedIndex=${task.lastProcessedTradeIndex}")
continue
}
// 记录当前处理的索引
lastProcessedIndexInPage = index
// 更新进度
val progress = if (pageTrades.size > 0) {
@@ -380,10 +393,10 @@ class BacktestExecutionService(
// 更新当前页的最后处理信息
val lastTradeInPage = currentPageTrades.lastOrNull()
if (lastTradeInPage != null) {
if (lastTradeInPage != null && lastProcessedIndexInPage != null) {
task.lastProcessedTradeTime = lastTradeInPage.tradeTime
task.lastProcessedTradeIndex = globalIndex + pageTrades.size - 1
task.processedTradeCount = task.lastProcessedTradeIndex + 1
task.lastProcessedTradeIndex = lastProcessedIndexInPage
task.processedTradeCount = lastProcessedIndexInPage + 1
task.finalBalance = currentBalance
backtestTaskRepository.save(task)
@@ -396,11 +409,7 @@ class BacktestExecutionService(
// 将当前页交易添加到全局列表(用于最终统计)
trades.addAll(currentPageTrades)
// 将当前页交易添加到全局列表(用于最终统计)
trades.addAll(currentPageTrades)
// 更新全局索引,准备处理下一页
globalIndex += pageTrades.size
// 准备处理下一页
currentPage++
} catch (e: Exception) {
@@ -93,25 +93,26 @@ class BacktestPollingService(
runBlocking {
// 支持恢复:如果有恢复点,计算从哪一页开始
val pageSize = 100
val page = if (currentTask.lastProcessedTradeIndex != null && currentTask.lastProcessedTradeIndex >= 0) {
// 从第几页开始(页码从 1 开始)
// 例如:已处理了99笔,lastProcessedTradeIndex=99,应从第2页开始
// 例如:已处理了0笔,lastProcessedTradeIndex=0,应从第1页开始(因为第1页的第1笔已经处理)
val lastProcessedIndex = currentTask.lastProcessedTradeIndex
val calculatedPage = (lastProcessedIndex / pageSize) + 1
val page = if (currentTask.lastProcessedTradeIndex != null) {
// 从第几页开始(页码从 0 开始)
// 例如:已处理了99笔,lastProcessedTradeIndex=99,应从第1页开始offset=100
val lastProcessedIndex = currentTask.lastProcessedTradeIndex!!
// 计算已处理的页码(从 0 开始)
val processedPage = lastProcessedIndex / pageSize
// 特殊情况:如果lastProcessedTradeIndex刚好是100的倍数减1(比如99,199,299...
// 说明该页已经完全处理,应该从下一页开始
val nextPage = if (lastProcessedIndex % pageSize == pageSize - 1) {
calculatedPage + 1
processedPage + 1
} else {
calculatedPage
processedPage
}
logger.info("恢复任务:已处理索引=$lastProcessedIndex, 计算页码=$nextPage, size=$pageSize")
nextPage
} else {
1 // 从第页开始
logger.info("新任务:从第0页开始")
0 // 从第0页开始(offset=0
}
logger.info("执行回测任务: taskId=${currentTask.id}, page=$page, size=$pageSize")
@@ -0,0 +1,15 @@
-- ============================================
-- 修复回测恢复逻辑:将 last_processed_trade_index 默认值改为 NULL
-- ============================================
-- 问题:新建任务的 last_processed_trade_index 默认值为 0,导致被误判为恢复任务
-- 解决:将默认值改为 NULL,并将现有新任务的 0 值改为 NULL
-- 1. 将现有新任务(status='PENDING' 且 last_processed_trade_index=0)的索引值改为 NULL
UPDATE backtest_task
SET last_processed_trade_index = NULL
WHERE status = 'PENDING' AND last_processed_trade_index = 0;
-- 2. 修改字段定义,允许 NULL 并设置默认值为 NULL
ALTER TABLE backtest_task
MODIFY COLUMN last_processed_trade_index INT DEFAULT NULL COMMENT '最后处理的交易索引(用于中断恢复)';