diff --git a/server-rs/crates/api-server/src/game_play_counter_worker.rs b/server-rs/crates/api-server/src/game_play_counter_worker.rs index 6e089007d..44d4148c8 100644 --- a/server-rs/crates/api-server/src/game_play_counter_worker.rs +++ b/server-rs/crates/api-server/src/game_play_counter_worker.rs @@ -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) { 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::(); + 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) { } } +fn total_delta(deltas: &[GamePlayCountDelta]) -> u64 { + deltas.iter().map(|delta| delta.delta).sum() +} + async fn write_batch( state: &AppState, deltas: &[GamePlayCountDelta],