From 75ad087e2cf3b6dcc27f4b9fe5a4facd01ad51da 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:07:02 +0800 Subject: [PATCH] =?UTF-8?q?=E5=90=8E=E7=AB=AF=EF=BC=9A=E5=85=B3=E5=81=9C?= =?UTF-8?q?=E4=B8=8D=E5=86=8D=E5=BC=BA=E5=88=B6=20flush=20=E6=B8=B8?= =?UTF-8?q?=E7=8E=A9=E8=AE=A1=E6=95=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - main.rs:finalize_shutdown 删除游戏游玩计数强制 flush,关停不再为它等待网络,内存里剩余增量随进程结束丢弃 - game_play_counter_worker.rs:删除 flush_game_play_counter_for_shutdown,flush_deltas 去掉 requeue_on_build 参数(关停路径已不存在),Build 仍放回等待下一轮 - game_play_counter.rs:take_pending 标记为仅测试使用,模块文档改指 take_pending_if_due --- .../api-server/src/game_play_counter.rs | 5 ++-- .../src/game_play_counter_worker.rs | 30 ++++--------------- server-rs/crates/api-server/src/main.rs | 20 ------------- 3 files changed, 9 insertions(+), 46 deletions(-) diff --git a/server-rs/crates/api-server/src/game_play_counter.rs b/server-rs/crates/api-server/src/game_play_counter.rs index 3eecc60c2..7fc5b8c46 100644 --- a/server-rs/crates/api-server/src/game_play_counter.rs +++ b/server-rs/crates/api-server/src/game_play_counter.rs @@ -3,7 +3,7 @@ //! 只在内存里累计「开始游戏」上报:按 `identity + game_id` 做去重窗口、按 `IP + game_id` 做固定 //! 窗口限流,并按 flush 间隔或待落库游戏数上限决定何时把增量交给写库方。 //! -//! 网络写入不在本模块内:`take_pending` / `requeue` 让调用方在锁外发起 procedure,锁只覆盖 +//! 网络写入不在本模块内:`take_pending_if_due` / `requeue` 让调用方在锁外发起 procedure,锁只覆盖 //! HashMap 操作。时间点全部由调用方传入 `now`,因此本模块不依赖运行时,可直接单测。 use std::{ @@ -179,7 +179,8 @@ impl GamePlayCounter { Some(sort_deltas(deltas)) } - /// 无条件取走全部增量,供关停 flush 使用。 + /// 无条件取走全部增量;仅供测试使用,生产路径走 `take_pending_if_due`。 + #[cfg(test)] pub fn take_pending(&self) -> Vec { let mut state = self.lock(); let deltas = drain_pending(&mut state.pending); 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 f740a1fab..6e089007d 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 @@ -3,6 +3,8 @@ //! 失败语义与 ADR 一致:连接还没建起来(`Build`)属于确定未发出,交给计数器放回下一轮重试; //! 其余错误(`Timeout` / `ConnectDropped` / `Procedure`)无法判断是否已提交,直接丢弃该批并 //! 记录丢失量,避免系统性双计。 +//! +//! 进程关停不做强制 flush:内存里剩下的增量随进程结束丢弃,关停路径不为它等待网络。 use std::time::{Duration, Instant}; @@ -28,32 +30,22 @@ pub(crate) fn spawn_game_play_counter_worker(state: AppState) { let counter = state.game_play_counter(); counter.prune_expired(now); if let Some(deltas) = counter.take_pending_if_due(now) { - flush_deltas(&state, deltas, true).await; + flush_deltas(&state, deltas).await; } } }); } -/// 关停前强制 flush 一次;由 `finalize_shutdown` 在总超时内调用。 -pub(crate) async fn flush_game_play_counter_for_shutdown(state: &AppState) { - let deltas = state.game_play_counter().take_pending(); - if deltas.is_empty() { - return; - } - flush_deltas(state, deltas, false).await; -} - /// 把一批增量写库。 /// -/// `requeue_on_build` 控制 `Build`(确定未发出)失败时是否放回:worker 循环为 `true`,等待下一轮; -/// 关停路径为 `false`——`finalize_shutdown` 之后不再重试,放回只会随进程退出丢失,应直接记丢失量。 -async fn flush_deltas(state: &AppState, deltas: Vec, requeue_on_build: bool) { +/// `Build`(确定未发出)放回等待下一轮;其余错误无法判断是否已提交,直接丢弃并记丢失量。 +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) { match write_batch(state, chunk).await { Ok(()) => accepted += chunk.len(), - Err(SpacetimeClientError::Build(message)) if requeue_on_build => { + Err(SpacetimeClientError::Build(message)) => { state.game_play_counter().requeue(chunk); warn!( error = %message, @@ -61,16 +53,6 @@ async fn flush_deltas(state: &AppState, deltas: Vec, requeue "游戏游玩计数写入连接未建立,已放回等待下一轮" ); } - Err(SpacetimeClientError::Build(message)) => { - let lost = chunk.iter().map(|delta| delta.delta).sum::(); - dropped = dropped.saturating_add(lost); - warn!( - error = %message, - games = chunk.len(), - lost, - "游戏游玩计数关停 flush 连接未建立,未落库增量丢弃" - ); - } Err(error) => { let lost = chunk.iter().map(|delta| delta.delta).sum::(); dropped = dropped.saturating_add(lost); diff --git a/server-rs/crates/api-server/src/main.rs b/server-rs/crates/api-server/src/main.rs index 4954845cc..6f270f9d0 100644 --- a/server-rs/crates/api-server/src/main.rs +++ b/server-rs/crates/api-server/src/main.rs @@ -682,26 +682,6 @@ async fn finalize_shutdown(context: ShutdownContext) { } } } - - if let Some(state) = context.app_state.as_ref() { - info!(timeout_ms, "api-server 退出前 flush 游戏游玩计数内存缓冲"); - match timeout( - context.outbox_flush_timeout, - crate::game_play_counter_worker::flush_game_play_counter_for_shutdown(state), - ) - .await - { - Ok(()) => { - info!("api-server 退出前游戏游玩计数 flush 完成"); - } - Err(_) => { - warn!( - timeout_ms, - "api-server 退出前游戏游玩计数 flush 超时,未落库增量已丢弃" - ); - } - } - } } fn spawn_common_app_state_background_workers(state: &AppState) {