去掉订阅快照的周期性重发
节点侧已改为按「收到主站的任何消息」判定新鲜度(过期 ⇔ 主站失联), 主站每 30 分钟的可用性检查即可持续刷新,因此主站不需要任何为续期的周期性推送。 - 删除上一版加入的 scheduledSubscriptionKeepalive(6 小时一次)。 - 保留「内容变化即推」(各处 requestSubscriptionSync)与「节点上线即推」 (initChannel),两者覆盖变更下发与节点冷启动两种场景。 - 移除对应的 standby 保活配置项。 最终推送时机:只在订阅内容变化时推送,另在节点上线时补推一次; 空闲时零快照传输(此前 60 秒一次时为约 486 MiB/天)。 测试:433 项全过。新增 SubscriptionSnapshotPushTest 锁住两条约定—— RemoteService 不得再有定时任务(续期改由节点存活信号承担),以及订阅刷新成功后 仍必须触发一次推送(用真实 refreshDueAccounts 调用验证)。
This commit is contained in:
@@ -18,7 +18,6 @@ import io.netty.util.concurrent.Promise;
|
||||
import lombok.Data;
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
import org.springframework.beans.factory.annotation.Value;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
import org.springframework.stereotype.Service;
|
||||
|
||||
import jakarta.annotation.PostConstruct;
|
||||
@@ -254,21 +253,10 @@ public class RemoteService {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 低频保活推送。
|
||||
*
|
||||
* <p>快照推送改为「订阅内容变化时立即推」之后,内容长期不变时节点将收不到任何推送,
|
||||
* 而节点的过期判定依据「最近一次收到快照的时刻」——于是它会在一周后永久判过期。
|
||||
* 这个保活用于把「主站在线」这件事周期性告诉节点,刷新其新鲜度。
|
||||
*
|
||||
* <p>周期远小于节点 `SubscriptionMaxStaleSeconds`(默认 7 天):默认 6 小时一次,
|
||||
* 约 1.4 MiB/天,相比原先每 60 秒一次的约 486 MiB/天 降低约三个数量级。
|
||||
*/
|
||||
@Scheduled(fixedDelayString = "${subscription.standby.keepalive-interval-ms:21600000}",
|
||||
initialDelayString = "${subscription.standby.keepalive-initial-delay-ms:600000}")
|
||||
void scheduledSubscriptionKeepalive() {
|
||||
requestSubscriptionSync();
|
||||
}
|
||||
// 这里刻意没有「定期重发整份快照」的定时任务:
|
||||
// 快照只在订阅内容变化时推送(各处 requestSubscriptionSync),节点上线时补推一次。
|
||||
// 备机的过期语义是「主站失联」,由节点按「收到主站的任何消息」判定新鲜度——
|
||||
// 主站每 30 分钟的可用性检查即可持续刷新它,因此不需要为续期而周期性传输快照。
|
||||
|
||||
private void drainSubscriptionSyncQueue() {
|
||||
try {
|
||||
|
||||
@@ -76,10 +76,9 @@ subscription:
|
||||
standby:
|
||||
sync-enabled: "${SUBSCRIPTION_STANDBY_SYNC_ENABLED:false}"
|
||||
sync-secret: "${SUBSCRIPTION_SYNC_SECRET:}"
|
||||
# 快照改为「订阅内容变化时立即推送」,另加低频保活刷新节点新鲜度。
|
||||
# 周期须远小于节点 SubscriptionMaxStaleSeconds(默认 7 天),默认 6 小时。
|
||||
keepalive-interval-ms: 21600000
|
||||
keepalive-initial-delay-ms: 600000
|
||||
# 快照只在订阅内容变化时推送(另在节点上线时补推一次),没有周期性重发。
|
||||
# 备机的过期语义是「主站失联」,由节点按「收到主站的任何消息」判定新鲜度,
|
||||
# 主站每 30 分钟的可用性检查即可持续刷新,无需为续期周期性传输快照。
|
||||
|
||||
bot:
|
||||
token: "5222939329:AAHa6l9ZuVVdNSDLPI_H-c8O_VgeOEw5plA"
|
||||
|
||||
@@ -1,52 +0,0 @@
|
||||
package com.lion.lionwebsite.Service;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* 备机快照推送时机与保活周期。
|
||||
*
|
||||
* <p>快照推送是「订阅内容变化时立即推」(由各处 requestSubscriptionSync 触发)。
|
||||
* 但节点按「最近一次收到快照的时刻」判过期,因此内容长期不变时必须有低频保活,
|
||||
* 否则节点会在一个有效期后永久判过期——这是本类要锁住的不变量。
|
||||
*/
|
||||
class SubscriptionKeepaliveTest {
|
||||
|
||||
/** 节点侧 SubscriptionMaxStaleSeconds 的默认值(7 天),保活周期必须远小于它。 */
|
||||
private static final long NODE_MAX_STALE_MILLIS = 7L * 24 * 3600 * 1000;
|
||||
|
||||
/**
|
||||
* 保活必须存在、且周期远小于节点有效期。
|
||||
*
|
||||
* <p>回归背景:曾一度改成「只在内容变化时推送」,结果存储节点重启后收不到任何推送,
|
||||
* 状态直接退化为 expired。保活就是为这个场景兜底的。
|
||||
*/
|
||||
@Test
|
||||
void keepaliveIsScheduledWellWithinNodeStaleWindow() throws Exception {
|
||||
Method method = RemoteService.class.getDeclaredMethod("scheduledSubscriptionKeepalive");
|
||||
Scheduled scheduled = method.getAnnotation(Scheduled.class);
|
||||
|
||||
assertNotNull(scheduled, "必须存在保活推送的定时入口");
|
||||
long interval = parseDefaultMillis(scheduled.fixedDelayString());
|
||||
assertTrue(interval > 0, "保活周期必须为正: " + scheduled.fixedDelayString());
|
||||
assertTrue(interval < NODE_MAX_STALE_MILLIS / 4,
|
||||
"保活周期(" + interval + "ms)必须远小于节点有效期(" + NODE_MAX_STALE_MILLIS + "ms),"
|
||||
+ "否则节点可能在两次保活之间判过期");
|
||||
|
||||
long initialDelay = parseDefaultMillis(scheduled.initialDelayString());
|
||||
assertTrue(initialDelay >= 0, "初始延迟不能为负");
|
||||
assertTrue(initialDelay < interval, "初始延迟应小于一个周期,避免启动后长时间不刷新节点新鲜度");
|
||||
}
|
||||
|
||||
/** 从 `${key:default}` 形式的占位符里取出默认毫秒值。 */
|
||||
private static long parseDefaultMillis(String placeholder) {
|
||||
int colon = placeholder.indexOf(':');
|
||||
String value = colon >= 0 ? placeholder.substring(colon + 1) : placeholder;
|
||||
value = value.replace("}", "").trim();
|
||||
return Long.parseLong(value);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
package com.lion.lionwebsite.Service;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.springframework.scheduling.annotation.Scheduled;
|
||||
|
||||
import java.lang.reflect.Method;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* 备机快照的推送时机。
|
||||
*
|
||||
* <p>约定:快照**只在订阅内容变化时**推送(各处 {@code requestSubscriptionSync}),
|
||||
* 另在节点上线时补推一次。这里刻意不做周期性重发——备机的过期语义是「主站失联」,
|
||||
* 由节点按「收到主站的任何消息」判定新鲜度(见 storageNode 侧 {@code markPrimaryContact}),
|
||||
* 主站每 30 分钟的可用性检查就能持续刷新它。
|
||||
*
|
||||
* <p>历史教训:曾用「每 60 秒重发整份快照」来续期,12 账号时约 346 KiB/次、
|
||||
* 约 486 MiB/天,而节点每次完整校验后都判 APPLY_OLD 丢弃;也曾改成「只在变化时推」
|
||||
* 却不给节点任何存活信号,导致内容长期不变时备机在第 7 天永久判过期。
|
||||
* 本类锁住「两者都不再发生」。
|
||||
*/
|
||||
class SubscriptionSnapshotPushTest {
|
||||
|
||||
/**
|
||||
* RemoteService 不得再有周期性重发快照的定时入口。
|
||||
*
|
||||
* <p>续期职责已下沉到节点侧(任何主站消息都刷新新鲜度),主站侧的周期重发是纯粹的浪费。
|
||||
*/
|
||||
@Test
|
||||
void hasNoPeriodicSnapshotResend() {
|
||||
for (Method method : RemoteService.class.getDeclaredMethods()) {
|
||||
Scheduled scheduled = method.getAnnotation(Scheduled.class);
|
||||
if (scheduled == null)
|
||||
continue;
|
||||
fail("RemoteService 不应再有定时任务,但发现: " + method.getName()
|
||||
+ "(订阅快照续期应由节点的存活信号判定承担)");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 订阅状态变更的入口必须仍然主动触发推送,否则变化无法在数秒内到达备机。
|
||||
*
|
||||
* <p>用真实调用验证:刷新调度器完成一轮刷新后,必须调用过一次 requestSubscriptionSync。
|
||||
*/
|
||||
@Test
|
||||
void subscriptionRefreshTriggersPush() {
|
||||
com.lion.lionwebsite.Dao.normal.SubMapper subMapper =
|
||||
org.mockito.Mockito.mock(com.lion.lionwebsite.Dao.normal.SubMapper.class);
|
||||
RemoteService remoteService = org.mockito.Mockito.mock(RemoteService.class);
|
||||
SubscriptionRefreshService refreshService =
|
||||
org.mockito.Mockito.mock(SubscriptionRefreshService.class);
|
||||
SubscriptionRefreshPlanner planner = new SubscriptionRefreshPlanner(
|
||||
java.time.Clock.systemDefaultZone(), new java.util.Random(),
|
||||
java.time.Duration.ofHours(24), java.time.Duration.ofHours(1), java.time.Duration.ofMinutes(5));
|
||||
SubscriptionRefreshScheduler scheduler = new SubscriptionRefreshScheduler(
|
||||
subMapper, refreshService, planner,
|
||||
org.mockito.Mockito.mock(PushService.class), remoteService);
|
||||
scheduler.maxPerTick = 2; // 手工构造的实例不受 @Value 注入,需显式设置单轮上限
|
||||
|
||||
com.lion.lionwebsite.Domain.SubscriptionAccount account =
|
||||
new com.lion.lionwebsite.Domain.SubscriptionAccount();
|
||||
account.setId(1);
|
||||
account.setName("acc");
|
||||
account.setUpstreamKey("key");
|
||||
account.setEnabled(true);
|
||||
account.setNextRefreshAt(System.currentTimeMillis() - 1); // 已到期
|
||||
account.setLastSuccessEpoch(System.currentTimeMillis());
|
||||
org.mockito.Mockito.when(subMapper.selectAllSubscriptionAccounts())
|
||||
.thenReturn(new java.util.ArrayList<>(java.util.List.of(account)));
|
||||
org.mockito.Mockito.when(refreshService.refresh(1)).thenReturn(true);
|
||||
|
||||
scheduler.refreshDueAccounts();
|
||||
|
||||
org.mockito.Mockito.verify(remoteService).requestSubscriptionSync();
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user