后端:游玩计数 flush 失败后直接丢弃剩余分片

- game_play_counter_worker.rs:任一分片失败即停止本次 flush,不再为后续分片逐个重建连接(连接不通时每片都会各自等一次超时)
- Build(确定未发出)仍放回本批,剩余分片直接丢弃;其余错误连本批一起丢弃,并记录丢弃增量
- 补 total_delta 汇总剩余增量,模块文档写明失败即停止与丢弃口径
This commit is contained in:
2026-10-04 11:10:03 +08:00
parent 75ad087e2c
commit 28ac8a1037
@@ -4,6 +4,9 @@
//! 其余错误(`Timeout` / `ConnectDropped` / `Procedure`)无法判断是否已提交,直接丢弃该批并
//! 记录丢失量,避免系统性双计。
//!
//! 任一分片失败都会终止本次 flush 的后续分片:连接不通时剩余分片只会重复同样的失败,逐个重试
//! 会把 worker 卡在多次连接超时上;剩余增量按“直接丢弃”处理,尽快回到 tick。
//!
//! 进程关停不做强制 flush:内存里剩下的增量随进程结束丢弃,关停路径不为它等待网络。
use std::time::{Duration, Instant};
@@ -42,26 +45,40 @@ pub(crate) fn spawn_game_play_counter_worker(state: AppState) {
async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>) {
let mut accepted = 0usize;
let mut dropped = 0u64;
for chunk in deltas.chunks(GAME_PLAY_COUNT_FLUSH_BATCH_SIZE) {
let mut start = 0usize;
while start < deltas.len() {
let end = (start + GAME_PLAY_COUNT_FLUSH_BATCH_SIZE).min(deltas.len());
let chunk = &deltas[start..end];
match write_batch(state, chunk).await {
Ok(()) => accepted += chunk.len(),
Ok(()) => {
accepted += chunk.len();
start = end;
}
// 连接没建起来(确定未发出):本批放回,但后面的分片会重复同样的失败,
// 直接丢弃剩余,避免每个分片各等一次连接超时把 worker 卡住。
Err(SpacetimeClientError::Build(message)) => {
state.game_play_counter().requeue(chunk);
let lost = total_delta(&deltas[end..]);
dropped = dropped.saturating_add(lost);
warn!(
error = %message,
games = chunk.len(),
"游戏游玩计数写入连接未建立,已放回等待下一轮"
lost,
"游戏游玩计数写入连接未建立:本批放回等待下一轮,剩余增量直接丢弃"
);
break;
}
// 结果未知:连本批一起丢弃剩余,避免双计,也不再把 worker 卡在逐个重连上。
Err(error) => {
let lost = chunk.iter().map(|delta| delta.delta).sum::<u64>();
let lost = total_delta(&deltas[start..]);
dropped = dropped.saturating_add(lost);
warn!(
error = %error,
games = chunk.len(),
games = deltas.len() - start,
lost,
"游戏游玩计数写入结果未知,丢弃本批以避免双计"
"游戏游玩计数写入失败,剩余增量直接丢弃"
);
break;
}
}
}
@@ -73,6 +90,10 @@ async fn flush_deltas(state: &AppState, deltas: Vec<GamePlayCountDelta>) {
}
}
fn total_delta(deltas: &[GamePlayCountDelta]) -> u64 {
deltas.iter().map(|delta| delta.delta).sum()
}
async fn write_batch(
state: &AppState,
deltas: &[GamePlayCountDelta],