using Microsoft.Extensions.Caching.Distributed; using Microsoft.Extensions.Logging; using Newtonsoft.Json.Linq; using PARR.Core.Common.Interfaces; using PARR.Infrastructure.Redis.Helpers; using StackExchange.Redis; using System.Text.Json; namespace PARR.Infrastructure.Redis { internal class RedisCacheService : IRedisCacheService { /// /// Минимальный размер для сжатия /// private const int compressionTreshold = 10 * 1024; // 10Kb private readonly IDistributedCache cache; private readonly IConnectionMultiplexer connectionMultiplexer; private readonly ILogger logger; private readonly IDatabase redis; public RedisCacheService( IDistributedCache cache, IConnectionMultiplexer connectionMultiplexer, ILogger logger ) { this.cache = cache; this.connectionMultiplexer = connectionMultiplexer; this.logger = logger; this.redis = connectionMultiplexer.GetDatabase(); } #region Распределенный кэш IDistributedCache public async Task GetCachedDataAsync(string key, bool useCompression = false) { if (useCompression) { // хранили в байтах try { var data = await cache.GetAsync(key); if (data == null) return default(T); return DeserializeWithCompression(data); } catch (Exception ex) { throw new InvalidOperationException("Формат данных не байты.", ex); } } else { // хранили в строке try { var jsonData = await cache.GetStringAsync(key); if (jsonData == null) return default(T); return JsonSerializer.Deserialize(jsonData); } catch (Exception ex) { throw new InvalidOperationException("Формат данных не строка.", ex); } } } //public T? GetCachedData(string key) //{ // var jsonData = cache.GetString(key); // if (jsonData == null) // return default(T); // return JsonSerializer.Deserialize(jsonData); //} //public void SetCachedData(string key, T data, TimeSpan cacheDuration) //{ // var options = new DistributedCacheEntryOptions // { // AbsoluteExpirationRelativeToNow = cacheDuration // }; // var jsonData = JsonSerializer.Serialize(data); // cache.SetString(key, jsonData, options); //} public async Task SetCachedDataAsync(string key, T data, TimeSpan cacheDuration, bool useCompression = false) { var options = new DistributedCacheEntryOptions { AbsoluteExpirationRelativeToNow = cacheDuration }; if (useCompression) { // храним в байтах, если надо то сжимаем var byteData = SerializeWithCompression(data); await cache.SetAsync(key, byteData, options); } else { // храним в строке var jsonData = JsonSerializer.Serialize(data); await cache.SetStringAsync(key, jsonData, options); } } public async Task DeleteCachedDataAsync(string key) { await cache.RemoveAsync(key); } public void DeleteCachedData(string key) { cache.Remove(key); } #endregion public async Task DeleteKeysByPatternAsync(string pattern, int batchSize = 100) { if (string.IsNullOrEmpty(pattern)) return 0; logger.LogDebug("Начало удаление ключей по маске '{Pattern}' из кэш", pattern); var deletedCount = 0; var keysToDelete = new List(); // Получаем сервер для выполнения SCAN //todo: Для кластера нужно обрабатывать каждый узел отдельно var server = connectionMultiplexer.GetEndPoints() .Select(t => connectionMultiplexer.GetServer(t)) .FirstOrDefault(); if (server == null) { logger.LogError("Не удалось получить сервер Redis для операции SCAN"); return 0; } // server.Keys() лениво сканирует ключи через SCAN (не блокирует!) // Важно: не вызывай .ToList() — это загрузит всё в память! foreach (var key in server.Keys(pattern: pattern, pageSize: batchSize)) { keysToDelete.Add(key); // Удаляем пакетом, чтобы не делать лишний сетевой запрос на каждый ключ if (keysToDelete.Count >= batchSize) { await redis.KeyDeleteAsync(keysToDelete.ToArray()); deletedCount += keysToDelete.Count; logger.LogDebug("Удалено {Count} ключей по масске '{Pattern}'", keysToDelete.Count, pattern); keysToDelete.Clear(); } } // Удаляем остаток, если остался if (keysToDelete.Count > 0) { await redis.KeyDeleteAsync(keysToDelete.ToArray()); deletedCount += keysToDelete.Count; } logger.LogInformation("Завершено удаление ключей по маске '{Pattern}'. Удалено: {DeletedCount}", pattern, deletedCount); return deletedCount; } #region Нативные операции Redis, Redis Hash public async Task SetHashFieldAsync(string hashKey, string field, T value, TimeSpan? ttl = null, bool useCompression = false) { // меняет одно поле в Hash if (useCompression) { // храним в байтах, если надо то сжимаем var byteData = SerializeWithCompression(value); // true - поля не было, создалось новое. false - поле было, обновили значение var result = await redis.HashSetAsync(hashKey, field, byteData); } else { // храним в строке var jsonData = JsonSerializer.Serialize(value); // true - поля не было, создалось новое. false - поле было, обновили значение var result = await redis.HashSetAsync(hashKey, field, jsonData); } // если указан ttl, обновим для всего Hash // если не указан и ранее был создан hashKey, оставит его ttl; а если hashKey не было, то создаст его БЕССРОЧНЫМ!!! if (ttl.HasValue) await SetHashTtlAsync(hashKey, ttl.Value); } public async Task GetHashFieldAsync(string hashKey, string field, bool useCompression = false) { // получить значение поля из Hash var value = await redis.HashGetAsync(hashKey, field); if (value.IsNullOrEmpty) return default; if (useCompression) { // хранили в байтах try { return DeserializeWithCompression(value); } catch (Exception ex) { throw new InvalidOperationException("Формат данных не байты.", ex); } } else { // хранили в строке return JsonSerializer.Deserialize(value); } } public async Task> GetHashFieldsAsync(List<(string HashKey, string FieldKey)> keys, bool useCompression = false) { var allItems = new List(keys.Count); // Делим на порции по 5000 (оптимально) foreach (var chunk in keys.Chunk(5000)) { // Формируем пайплайн ТОЛЬКО для 5 000 элементов var tasks = new List>(chunk.Length); foreach (var key in chunk) { tasks.Add(redis.HashGetAsync(key.HashKey, key.FieldKey)); } logger.LogDebug("Сформировал {Count} одновременных запросов в кэш", tasks.Count); // Ждем ответа от текущей порции RedisValue[] results = await Task.WhenAll(tasks); // Парсим foreach (var result in results) { if (result.HasValue) { if (useCompression) { try { var item = DeserializeWithCompression((byte[])result); if (item != null) allItems.Add(item); } catch (Exception ex) { throw new InvalidOperationException("Формат данных не байты.", ex); } } else { var item = JsonSerializer.Deserialize(result.ToString()); if (item != null) allItems.Add(item); } } } } logger.LogDebug("Получил {Count} элементов из кэш.", allItems.Count); return allItems; } public async Task<(string Field, T? Value)?> GetFirstHashFieldAsync(string hashKey, bool useCompression = false) { logger.LogDebug("Запрос первого поля из хеша '{HashKey}'", hashKey); int attempts = 0; int maxAttempts = 5; await foreach (var entry in redis.HashScanAsync(hashKey, pageSize: 1)) { if (++attempts > maxAttempts) { logger.LogWarning("Превышен лимит попыток ({Max}) для хеша '{HashKey}'", maxAttempts, hashKey); return null; } if (entry.Value.IsNullOrEmpty) { logger.LogDebug("Поле '{Field}' в хеше '{HashKey}' пустое, пропускаем", entry.Name, hashKey); continue; } T? deserialized = default; if (useCompression) { // хранили в байтах try { deserialized = DeserializeWithCompression(entry.Value); } catch (Exception ex) { throw new InvalidOperationException("Формат данных не байты.", ex); } } else { deserialized = JsonSerializer.Deserialize(entry.Value); } logger.LogDebug("Получено поле '{Field}' из хеша '{HashKey}', тип: {Type}", entry.Name, hashKey, typeof(T).Name); return (entry.Name.ToString(), deserialized); } logger.LogDebug("Хеш '{HashKey}' пуст или не существует", hashKey); return null; } public async Task> GetAllHashFieldsAsync(string hashKey, bool useCompression = false) { // получить все записи из Hash var objs = await redis.HashGetAllAsync(hashKey); if (objs.Length == 0) return new Dictionary(); try { var result = objs.ToDictionary( t => t.Name.ToString(), t => useCompression ? DeserializeWithCompression(t.Value) : JsonSerializer.Deserialize(t.Value) ); return result!; } catch (Exception ex) { throw new Exception("Не смог десериализовать полученные данные", ex); //return new Dictionary(); } } public async Task DeleteHashFieldAsync(string hashKey, string field) { // удалить запись из Hash await redis.HashDeleteAsync(hashKey, field); } public async Task HashFieldExistsAsync(string hashKey, string field) { // есть ли запись в Hash return await redis.HashExistsAsync(hashKey, field); } public async Task DeleteHashAsync(string hashKey) { // удалить весь Hash await redis.KeyDeleteAsync(hashKey); } public async Task GetHashLengthAsync(string hashKey) { // кол-во записей в hash return await redis.HashLengthAsync(hashKey); } public async Task SetHashTtlAsync(string hashKey, TimeSpan ttl) { // Установить ttl для Hash await redis.KeyExpireAsync(hashKey, ttl); } #endregion #region Helpers public string GetKey(string[] keyParts, string[]? keyPartsToHash = null) { // разделитель между ИД //var mainSeparator = "_"; var mainSeparator = ":"; // разделитель между словами var wordSeparator = "-"; if (keyParts.Length == 0) { throw new ArgumentNullException("keyParts не может быть пустым"); } var keyStr = string.Join(mainSeparator, keyParts).Replace(" ", wordSeparator); if (keyPartsToHash != null && keyPartsToHash.Length > 0) { var partsToHashStr = string.Join(mainSeparator, keyPartsToHash).Replace(" ", wordSeparator); using var sha256 = System.Security.Cryptography.SHA256.Create(); var hashedBytes = sha256.ComputeHash(System.Text.Encoding.UTF8.GetBytes(partsToHashStr)); var base64str = Convert.ToBase64String(hashedBytes).Substring(0, 16); return (keyStr + mainSeparator + base64str).ToLower(); } return keyStr.ToLower(); } #endregion private byte[] SerializeWithCompression(T data) { var json = JsonSerializer.SerializeToUtf8Bytes(data); logger.LogDebug("Сериализовано: {Size} байт, тип: {Type}", json.Length, typeof(T).Name); if (json.Length < compressionTreshold) return json; var compressed = CompressionHelper.Compress(json); logger.LogDebug("После сжатия: {Size} байт (экономия: {Percent}%)", compressed.Length, Math.Round((1 - (double)compressed.Length / json.Length) * 100, 1)); return compressed; } private T? DeserializeWithCompression(byte[] data) { var json = CompressionHelper.Decompress(data); return JsonSerializer.Deserialize(json); } } }