using Microsoft.Extensions.Caching.Distributed;
using Microsoft.Extensions.Logging;
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
{
byte[] rawBytes = value!;
return DeserializeWithCompression(rawBytes);
}
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)
continue;
if (useCompression)
{
// Сжатые данные
try
{
byte[] rawBytes = result!;
var data = DeserializeWithCompression(rawBytes);
if (data != null)
allItems.Add(data!);
}
catch (Exception ex)
{
throw new InvalidOperationException("Формат данных не байты.", ex);
}
}
else
{
// Json
var data = JsonSerializer.Deserialize(result!.ToString());
if (data != null)
allItems.Add(data);
}
}
}
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 DeleteHashFieldsAsync(List<(string HashKey, string FieldKey)> keys)
{
if (keys == null || keys.Count == 0)
return;
int deletedCount = 0;
// Делим на порции по 5000, чтобы не перегружать буфер команд Redis
foreach (var chunk in keys.Chunk(5000))
{
var tasks = new List>(chunk.Length);
foreach (var key in chunk)
tasks.Add(redis.HashDeleteAsync(key.HashKey, key.FieldKey));
logger.LogDebug("Отправил пачку из {Count} команд на удаление в Redis", tasks.Count);
// Ждем выполнения текущей порции команд удалений
bool[] results = await Task.WhenAll(tasks);
// Считаем, сколько элементов реально было удалено
deletedCount += results.Count(x => x == true);
}
logger.LogDebug("Успешно удалено элементов из кэша: {Count}.", deletedCount);
}
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);
}
}
}