后端:关停不再强制 flush 游玩计数

- 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
This commit is contained in:
2026-10-04 11:07:02 +08:00
parent 3f0a77e65e
commit 75ad087e2c
3 changed files with 9 additions and 46 deletions
@@ -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<GamePlayCountDelta> {
let mut state = self.lock();
let deltas = drain_pending(&mut state.pending);
@@ -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<GamePlayCountDelta>, requeue_on_build: bool) {
/// `Build`(确定未发出)放回等待下一轮;其余错误无法判断是否已提交,直接丢弃并记丢失量。
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) {
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<GamePlayCountDelta>, requeue
"游戏游玩计数写入连接未建立,已放回等待下一轮"
);
}
Err(SpacetimeClientError::Build(message)) => {
let lost = chunk.iter().map(|delta| delta.delta).sum::<u64>();
dropped = dropped.saturating_add(lost);
warn!(
error = %message,
games = chunk.len(),
lost,
"游戏游玩计数关停 flush 连接未建立,未落库增量丢弃"
);
}
Err(error) => {
let lost = chunk.iter().map(|delta| delta.delta).sum::<u64>();
dropped = dropped.saturating_add(lost);
-20
View File
@@ -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) {