using System.Net.Http;
using System.Text.Json.Nodes;
using Dpz.Core.Public.Entity.Steam;
using Dpz.Core.Public.ViewModel.Response.Steam;
using Microsoft.AspNetCore.WebUtilities;
namespace Dpz.Core.Service.RepositoryServiceImpl;
public class SteamGameService(
IRepository<SteamGame> repository,
IMapper mapper,
HttpClient httpClient,
IConfiguration configuration,
ILogger<SteamGameService> logger,
IFusionCache fusionCache
) : ISteamGameService
{
private readonly Lazy<string> _lazyCdnHost = new(() =>
{
var cdnHost = configuration["S3:CdnHost"];
return string.IsNullOrWhiteSpace(cdnHost)
? throw new InvalidConfigurationException("S3:CdnHost is empty or null")
: cdnHost;
});
private readonly Lazy<string> _lazyLogoUrlTemplate = new(() =>
"https://shared.akamai.steamstatic.com/store_item_assets/steam/apps/{0}/header.jpg"
);
private readonly Lazy<string> _lazySteamId = new(() =>
{
var steamId = configuration["Steam:SteamId"];
return string.IsNullOrWhiteSpace(steamId)
? throw new InvalidConfigurationException("Steam:SteamId is empty or null")
: steamId;
});
private readonly Lazy<string> _lazySteamAppKey = new(() =>
{
var appKey = configuration["Steam:AppKey"];
return string.IsNullOrWhiteSpace(appKey)
? throw new InvalidConfigurationException("Steam:SteamAppKey is empty or null")
: appKey;
});
private readonly Lazy<List<int>> _filterIgnoreGameIds = new(() =>
configuration.GetSection("Steam:IgnoreFilterIds").Get<List<int>>() ?? []
);
// API请求延迟时间,防止请求过于频繁
private const int ApiDelayMs = 500;
/// <summary>
/// 当前实例ID,用于区分不同的实例
/// </summary>
private readonly Lazy<string> _instanceId = new(() => Guid.NewGuid().ToString("D"));
/// <summary>
/// 当前实例ID的缓存键
/// </summary>
private const string RunningInstanceIdKey = "SteamGameService.Running.InstanceId";
public event LogoDownload? OnLogoDownloadComplete;
public event AchievementIconGrayDownload? OnAchievementIconGrayDownloadComplete;
public event AchievementIconDownload? OnAchievementIconDownloadComplete;
[InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
public async Task UpdateGamesAsync(CancellationToken cancellationToken = default)
{
var runningInstanceId = await fusionCache.TryGetAsync<string>(
RunningInstanceIdKey,
token: cancellationToken
);
if (runningInstanceId.HasValue && runningInstanceId.Value != _instanceId.Value)
{
logger.LogInformation(
"当前已有实例正在运行,跳过此次更新,实例ID:{RunningInstanceId},当前实例ID:{CurrentInstanceId}",
runningInstanceId.Value,
_instanceId.Value
);
return;
}
await fusionCache.SetAsync(
RunningInstanceIdKey,
_instanceId.Value,
TimeSpan.FromHours(3),
token: cancellationToken
);
try
{
logger.LogInformation("开始更新Steam游戏数据,实例ID:{InstanceId}", _instanceId.Value);
var dbGames = await repository.SearchFor(x => true).ToListAsync(cancellationToken);
var games = await FetchGamesAsync(dbGames, cancellationToken);
// 安全检查:如果获取到空列表,直接返回,避免误删数据
if (games.Count == 0)
{
logger.LogWarning("获取到的Steam游戏数量为0,跳过此次更新以防止误删数据");
return;
}
// 新增游戏
var newGames = games.ExceptBy(dbGames.Select(x => x.Id), x => x.Id).ToList();
if (newGames.Count > 0)
{
logger.LogInformation("发现{Count}个新游戏,准备添加", newGames.Count);
await repository.InsertAsync(newGames, cancellationToken);
// var batchSize =
// newGames.Count % 10 > 0 ? newGames.Count / 10 + 1 : newGames.Count / 10;
//
// for (var i = 0; i < batchSize; i++)
// {
// var startIndex = i * 10;
// var count = Math.Min(10, newGames.Count - startIndex);
// var batch = newGames.GetRange(startIndex, count);
// logger.LogInformation(
// "批量添加游戏,开始索引:{StartIndex},数量:{Count}",
// startIndex,
// count
// );
// await repository.InsertAsync(batch);
// await Task.Delay(1000);
// }
}
// DB已有游戏,待更新
var updateGames = games.IntersectBy(dbGames.Select(x => x.Id), x => x.Id).ToList();
if (updateGames.Count > 0)
{
logger.LogInformation("开始更新{Count}个已有游戏", updateGames.Count);
await Parallel.ForEachAsync(
updateGames,
cancellationToken,
async (game, _) => await UpdateGameAsync(dbGames, game)
);
}
// 删除已删除的游戏
var deletedGames = dbGames.ExceptBy(games.Select(x => x.Id), x => x.Id).ToList();
if (deletedGames.Count > 0)
{
logger.LogInformation("发现{Count}个已删除游戏,准备删除", deletedGames.Count);
var filter = Builders<SteamGame>.Filter.In(
x => x.Id,
deletedGames.Select(x => x.Id)
);
await repository.DeleteAsync(filter, cancellationToken);
}
logger.LogInformation(
"Steam游戏数据更新完成,新增:{NewCount},更新:{UpdateCount},删除:{DeletedCount},实例ID:{InstanceId}",
newGames.Count,
updateGames.Count,
deletedGames.Count,
_instanceId.Value
);
}
catch (Exception ex)
{
logger.LogError(
ex,
"更新Steam游戏数据时发生异常,实例ID:{InstanceId}",
_instanceId.Value
);
throw;
}
finally
{
// 确保在任何情况下都释放锁
try
{
await fusionCache.RemoveAsync(RunningInstanceIdKey, token: cancellationToken);
logger.LogDebug("已释放Steam游戏更新锁,实例ID:{InstanceId}", _instanceId.Value);
}
catch (Exception ex)
{
logger.LogError(
ex,
"释放Steam游戏更新锁时发生异常,实例ID:{InstanceId}",
_instanceId.Value
);
}
}
}
private async Task UpdateGameAsync(List<SteamGame> dbGames, SteamGame updateGame)
{
var dbGame = dbGames.Find(x => x.Id == updateGame.Id);
if (dbGame != null)
{
var dic = dbGame.UpdateContent(updateGame);
if (dic.Count > 0)
{
logger.LogInformation(
"更新游戏 {GameId}/{GameName} 的数据,更新细节:{@UpdateContent}",
updateGame.Id,
updateGame.Name,
dic
);
}
}
await repository.UpdateAsync(updateGame);
}
private async Task<List<AchievementDetail>> FetchAchievementsAsync(
int id,
List<SteamGame> dbGames,
CancellationToken cancellationToken = default
)
{
var queryString = new Dictionary<string, string?>
{
{ "steamid", _lazySteamId.Value },
{ "key", _lazySteamAppKey.Value },
{ "format", "json" },
{ "appid", id.ToString() },
{ "l", "schinese" },
};
// 添加API请求延迟
await Task.Delay(ApiDelayMs, cancellationToken);
logger.LogDebug("获取游戏 {GameId} 的成就数据", id);
var achievements =
await ExecuteHttpRequestAsync<List<AchievementDetail>>(
"/ISteamUserStats/GetSchemaForGame/v2/",
queryString,
root => root?["game"]?["availableGameStats"]?["achievements"]
) ?? [];
var unlockAchievements = await FetchUnlockAchievementAsync(queryString);
#if DEBUG
var options = new ParallelOptions { MaxDegreeOfParallelism = 1 };
#endif
await Parallel.ForEachAsync(
achievements,
#if DEBUG
options,
#endif
async (achievement, _) =>
{
if (
achievement.Name != null
&& unlockAchievements.TryGetValue(achievement.Name, out var unlockTime)
&& unlockTime > 0
)
{
achievement.UnlockTime = unlockTime.ToDateTime();
}
await DownloadAchievementAsync(dbGames, id, achievement);
await DownloadAchievementGrayAsync(dbGames, id, achievement);
}
);
logger.LogDebug("游戏 {GameId} 共获取到 {Count} 个成就", id, achievements.Count);
return achievements;
}
private async Task<Dictionary<string, long>> FetchUnlockAchievementAsync(
Dictionary<string, string?> queryString
)
{
logger.LogDebug("获取已解锁的成就数据");
// 添加API请求延迟
await Task.Delay(ApiDelayMs);
return await ExecuteHttpRequestAsync<Dictionary<string, long>>(
"/ISteamUserStats/GetPlayerAchievements/v0001/",
queryString,
root =>
{
var achievements = root
?["playerstats"]?["achievements"]?.AsArray()
.Select(x =>
{
var name = x?["apiname"]?.GetValue<string>();
if (name == null)
{
logger.LogError(
"Achievement name is null for {SteamId}",
_lazySteamId.Value
);
return (KeyValuePair<string, long>?)null;
}
var unlockTime = x?["unlocktime"]?.GetValue<long>() ?? 0;
return x?["achieved"]?.GetValue<int>() == 1
? KeyValuePair.Create(name, unlockTime)
: KeyValuePair.Create(name, 0L);
})
.Where(x => x != null)
.Select(x => x!.Value)
.ToDictionary(x => x.Key, x => x.Value);
return JsonValue.Create(achievements);
}
) ?? [];
}
private async Task<List<SteamGame>> FetchGamesAsync(
List<SteamGame> dbGames,
CancellationToken cancellationToken = default
)
{
logger.LogInformation("开始获取Steam游戏列表");
var games = await ApplicationTools.RetryAsync(
async () =>
{
var queryString = new Dictionary<string, string?>
{
{ "steamid", _lazySteamId.Value },
{ "key", _lazySteamAppKey.Value },
{ "include_appinfo", "1" },
{ "format", "json" },
};
// 添加API请求延迟
await Task.Delay(ApiDelayMs, cancellationToken);
var result = await ExecuteHttpRequestAsync<List<SteamGame>>(
"/IPlayerService/GetOwnedGames/v0001/",
queryString,
root => root?["response"]?["games"]?.AsArray()
);
// 如果返回null,抛出异常以触发重试
if (result == null)
{
throw new InvalidOperationException("Steam API返回null,触发重试");
}
return result
// 已失效的游戏,导致后续请求图标、成就失败
// 不知道为什么steam不移除它
// 屏蔽
//.Where(x => x.Id != 692850)
.Where(x => !_filterIgnoreGameIds.Value.Contains(x.Id))
.ToList();
},
retryInterval: TimeSpan.FromMilliseconds(ApiDelayMs),
maxAttemptCount: 3,
specificExceptionTypes: [typeof(JsonException)] // JSON错误不重试
);
logger.LogInformation("获取到 {Count} 个Steam游戏", games.Count);
#if DEBUG
var options = new ParallelOptions { MaxDegreeOfParallelism = 1 };
#endif
await Parallel.ForEachAsync(
games,
#if DEBUG
options,
#endif
async (game, ct) => await UpdateGame(dbGames, game, ct)
);
return games;
}
private async Task UpdateGame(
List<SteamGame> dbGames,
SteamGame game,
CancellationToken cancellationToken = default
)
{
logger.LogDebug("更新游戏 {GameId}/{GameName} 的基本信息", game.Id, game.Name);
var achievements = await FetchAchievementsAsync(game.Id, dbGames, cancellationToken);
var originGame = dbGames.Find(x => x.Id == game.Id);
if (achievements is { Count: > 0 })
{
game.Achievements = achievements;
game.LastUpdateTime = DateTime.Now;
}
else if (originGame != null)
{
game.Achievements = originGame.Achievements;
}
else
{
game.Achievements = [];
game.LastUpdateTime = DateTime.Now;
}
if (!AreAssetsDownloaded(originGame))
{
logger.LogDebug("下载游戏 {GameId} 的Logo", game.Id);
// 添加API请求延迟
await Task.Delay(ApiDelayMs, cancellationToken);
await using var stream = await HttpDownloadAsync(
string.Format(_lazyLogoUrlTemplate.Value, game.Id)
);
if (stream != null && OnLogoDownloadComplete != null)
{
var url = await OnLogoDownloadComplete(stream, game.Id);
if (!string.IsNullOrWhiteSpace(url))
{
game.ImageLogo = url;
game.ImageIcon = url;
game.LastUpdateTime = DateTime.Now;
}
}
}
else
{
game.ImageLogo = originGame?.ImageLogo;
game.ImageIcon = originGame?.ImageIcon;
logger.LogDebug("游戏 {GameId} Logo 无需更新", game.Id);
}
}
private async Task<Stream?> HttpDownloadAsync(string uri)
{
try
{
var response = await httpClient.GetAsync(uri);
response.EnsureSuccessStatusCode();
var memoryStream = new MemoryStream();
await response.Content.CopyToAsync(memoryStream);
memoryStream.Position = 0;
return memoryStream;
}
catch (Exception e)
{
logger.LogError(e, "Download failed for {Uri}", uri);
return null;
}
}
private bool AreAssetsDownloaded(SteamGame? originGame)
{
return originGame != null
&& originGame.ImageIcon?.StartsWith(
_lazyCdnHost.Value,
StringComparison.OrdinalIgnoreCase
) == true
&& originGame.ImageLogo?.StartsWith(
_lazyCdnHost.Value,
StringComparison.OrdinalIgnoreCase
) == true;
}
private async Task DownloadAchievementAsync(
List<SteamGame> dbGames,
int id,
AchievementDetail achievement
)
{
var dbAchievement = dbGames
.FirstOrDefault(x => x.Id == id)
?.Achievements.Find(x => x.Name == achievement.Name);
if (
dbAchievement?.Icon?.StartsWith(_lazyCdnHost.Value, StringComparison.OrdinalIgnoreCase)
== true
)
{
achievement.Icon = dbAchievement.Icon;
return;
}
if (
string.IsNullOrWhiteSpace(achievement.Icon)
|| string.IsNullOrWhiteSpace(achievement.Name)
)
{
return;
}
logger.LogDebug("下载游戏 {GameId} 成就 {AchievementName} 的图标", id, achievement.Name);
// 添加API请求延迟
await Task.Delay(ApiDelayMs);
await using var icon = await HttpDownloadAsync(achievement.Icon);
var handler = OnAchievementIconDownloadComplete;
if (icon != null && handler != null)
{
var url = await handler.Invoke(icon, id, achievement.Name);
if (!string.IsNullOrWhiteSpace(url))
{
achievement.Icon = url;
}
}
}
private async Task DownloadAchievementGrayAsync(
IEnumerable<SteamGame> dbGames,
int id,
AchievementDetail achievement
)
{
var dbAchievement = dbGames
.FirstOrDefault(x => x.Id == id)
?.Achievements.Find(x => x.Name == achievement.Name);
if (
dbAchievement?.IconGray?.StartsWith(
_lazyCdnHost.Value,
StringComparison.OrdinalIgnoreCase
) == true
)
{
achievement.IconGray = dbAchievement.IconGray;
return;
}
if (
string.IsNullOrWhiteSpace(achievement.IconGray)
|| string.IsNullOrWhiteSpace(achievement.Name)
)
{
return;
}
logger.LogDebug(
"下载游戏 {GameId} 成就 {AchievementName} 的灰色图标",
id,
achievement.Name
);
// 添加API请求延迟
await Task.Delay(ApiDelayMs);
await using var icon = await HttpDownloadAsync(achievement.IconGray);
var handler = OnAchievementIconGrayDownloadComplete;
if (icon != null && handler != null)
{
var url = await handler.Invoke(icon, id, achievement.Name);
if (!string.IsNullOrWhiteSpace(url))
{
achievement.IconGray = url;
}
}
}
[Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
public async Task<List<SteamGameResponse>> GetGamesAsync(
CancellationToken cancellationToken = default
)
{
var list = await repository.SearchFor(x => true).ToListAsync(cancellationToken);
var result = mapper.Map<List<SteamGameResponse>>(list);
return result;
}
[Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
public async Task<SteamGameResponse?> GetGameAsync(
int id,
CancellationToken cancellationToken = default
)
{
var game = await repository
.SearchFor(x => x.Id == id)
.FirstOrDefaultAsync(cancellationToken);
if (game == null)
{
logger.LogWarning("未找到游戏 {GameId} 的信息", id);
return null;
}
var result = mapper.Map<SteamGameResponse>(game);
return result;
}
[InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
public Task ClearCacheAsync(CancellationToken cancellationToken = default)
{
return Task.CompletedTask;
}
private async Task<T?> ExecuteHttpRequestAsync<T>(
string endpoint,
Dictionary<string, string?> queryParams,
Func<JsonNode?, JsonNode?>? dataSelector = null
)
{
var requestUri = QueryHelpers.AddQueryString(endpoint, queryParams);
var jsonContent = "";
try
{
var response = await httpClient.GetAsync(requestUri);
response.EnsureSuccessStatusCode();
jsonContent = await response.Content.ReadAsStringAsync();
if (string.IsNullOrWhiteSpace(jsonContent))
{
logger.LogWarning("Steam API返回空内容,请求地址:{RequestUri}", requestUri);
return default;
}
var jsonNode = JsonNode.Parse(jsonContent);
if (jsonNode == null)
{
logger.LogWarning(
"Steam API返回的JSON无法解析,请求地址:{RequestUri},内容:{Content}",
requestUri,
jsonContent
);
return default;
}
var data = dataSelector?.Invoke(jsonNode) ?? jsonNode;
var result = data.Deserialize<T>();
if (result == null)
{
logger.LogWarning(
"Steam API数据反序列化失败,请求地址:{RequestUri},数据:{Data}",
requestUri,
data.ToString()
);
}
return result;
}
catch (JsonException e)
{
logger.LogError(
e,
"Steam API数据解析失败,请求地址:{RequestUri},内容:{Content}",
requestUri,
jsonContent
);
return default;
}
catch (Exception e)
{
logger.LogError(e, "Steam API请求失败,请求地址:{RequestUri}", requestUri);
return default;
}
}
}
⚠⚠ 以下内容为AI分析的结果,请根据实际情况进行判断。
代码解释
这是一个用于管理和同步 Steam 游戏数据的服务类 SteamGameService,主要功能包括从 Steam API 获取游戏信息、成就数据,以及下载和管理游戏相关图片资源。
核心功能模块
1. 依赖注入与配置初始化
public SteamGameService(
IRepository<SteamGame> repository,
IMapper mapper,
HttpClient httpClient,
IConfiguration configuration,
ILogger<SteamGameService> logger,
IFusionCache fusionCache
)
- 使用主构造函数注入依赖项
- 通过
Lazy<T>延迟加载配置项(CDN 主机、Steam ID、API Key 等) - 配置项为空时抛出
InvalidConfigurationException
2. 游戏数据更新(核心方法)
UpdateGamesAsync 方法实现了完整的数据同步流程:
分布式锁机制
var runningInstanceId = await fusionCache.TryGetAsync<string>(RunningInstanceIdKey, ...);
- 使用
FusionCache实现分布式锁,防止多实例并发执行 - 每个实例生成唯一 ID,锁有效期 3 小时
finally块确保锁一定被释放
数据同步逻辑
- 新增游戏:从 API 获取的游戏列表中排除数据库已有的
- 更新游戏:对数据库中已存在的游戏进行增量更新
- 删除游戏:移除 API 中不存在但数据库有的记录
- 安全检查:API 返回空列表时跳过更新,防止误删
3. Steam API 交互
API 请求频率控制
private const int ApiDelayMs = 500;
await Task.Delay(ApiDelayMs, cancellationToken);
- 每次请求前延迟 500ms,避免触发 Steam API 限流
重试机制
var games = await ApplicationTools.RetryAsync(
async () => { ... },
retryInterval: TimeSpan.FromMilliseconds(ApiDelayMs),
maxAttemptCount: 3,
specificExceptionTypes: [typeof(JsonException)]
);
- 最多重试 3 次
- JSON 解析错误不重试(数据格式问题)
通用 HTTP 请求封装
ExecuteHttpRequestAsync<T> 方法:
- 支持自定义 JSON 数据提取器 (
dataSelector) - 统一异常处理和日志记录
- 返回泛型结果,增强复用性
4. 成就数据处理
获取成就详情
private async Task<List<AchievementDetail>> FetchAchievementsAsync(...)
- 从
/ISteamUserStats/GetSchemaForGame/v2/获取成就列表 - 从
/ISteamUserStats/GetPlayerAchievements/v0001/获取解锁状态 - 合并两个接口数据,填充
UnlockTime字段
并行下载成就图标
await Parallel.ForEachAsync(achievements, async (achievement, _) => {
await DownloadAchievementAsync(...);
await DownloadAchievementGrayAsync(...);
});
- 同时下载彩色和灰色图标
- Debug 模式下限制并行度为 1(方便调试)
5. 资源下载与事件通知
委托事件机制
public event LogoDownload? OnLogoDownloadComplete;
public event AchievementIconGrayDownload? OnAchievementIconGrayDownloadComplete;
public event AchievementIconDownload? OnAchievementIconDownloadComplete;
- 下载完成后触发事件,由外部处理文件上传到 S3
- 返回 CDN URL 后更新到实体对象
下载优化
private bool AreAssetsDownloaded(SteamGame? originGame)
{
return originGame?.ImageIcon?.StartsWith(_lazyCdnHost.Value, ...) == true;
}
- 检查资源是否已存在于 CDN,避免重复下载
6. 缓存策略
方法级缓存
[Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
public async Task<List<SteamGameResponse>> GetGamesAsync(...)
- 查询结果缓存 1 天
缓存失效
[InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
public async Task UpdateGamesAsync(...)
- 更新数据时自动清除相关缓存
- 通过 AOP 实现(需要配合拦截器框架)
7. 增量更新机制
private async Task UpdateGameAsync(List<SteamGame> dbGames, SteamGame updateGame)
{
var dic = dbGame.UpdateContent(updateGame);
if (dic.Count > 0) {
logger.LogInformation("更新细节:{@UpdateContent}", dic);
}
}
- 只更新变化的字段
- 记录详细的变更内容用于审计
设计亮点
- 防御性编程:配置校验、空值检查、异常重试
- 分布式友好:实例锁、幂等性设计
- 性能优化:并行处理、延迟加载、智能缓存
- 可观测性:详细的日志记录(包含实例 ID、游戏 ID 等上下文)
- 扩展性:事件驱动的资源处理机制
潜在改进点
Parallel.ForEachAsync可配置并发度(当前硬编码)- API 延迟时间可配置化
- 批量更新的注释代码可能需要恢复(大量新游戏场景)
- 错误重试策略可进一步细化(如指数退避)
AI 正在分析代码…
评论加载中...