From 28ac8a1037b92d012ee7629d9cccfcc4372514cd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Sun, 4 Oct 2026 11:10:03 +0800 Subject: [PATCH] =?UTF-8?q?=E5=90=8E=E7=AB=AF=EF=BC=9A=E6=B8=B8=E7=8E=A9?= =?UTF-8?q?=E8=AE=A1=E6=95=B0=20flush=20=E5=A4=B1=E8=B4=A5=E5=90=8E?= =?UTF-8?q?=E7=9B=B4=E6=8E=A5=E4=B8=A2=E5=BC=83=E5=89=A9=E4=BD=99=E5=88=86?= =?UTF-8?q?=E7=89=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - game_play_counter_worker.rs:任一分片失败即停止本次 flush,不再为后续分片逐个重建连接(连接不通时每片都会各自等一次超时) - Build(确定未发出)仍放回本批,剩余分片直接丢弃;其余错误连本批一起丢弃,并记录丢弃增量 - 补 total_delta 汇总剩余增量,模块文档写明失败即停止与丢弃口径 --- .../src/game_play_counter_worker.rs | 33 +++++++++++++++---- 1 file changed, 27 insertions(+), 6 deletions(-) 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],