safe_redis_cli.py 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  1. import time
  2. import redis
  3. from typing import List, Dict, Callable, Any
  4. class SafeRedisClient:
  5. def __init__(self, redis_conf: Dict, logger: Callable[[str], None] = None):
  6. safe_conf = dict(redis_conf)
  7. safe_conf.setdefault('socket_timeout', 5.0) # 读写超时防死等
  8. safe_conf.setdefault('socket_connect_timeout', 5.0) # 建连超时防死等
  9. safe_conf.setdefault('socket_keepalive', True) # TCP 操作系统级保活
  10. safe_conf.setdefault('health_check_interval', 30) # 自动探测断连并重连
  11. self._client = redis.Redis(**safe_conf)
  12. self._logger = logger
  13. def _log(self, msg: str):
  14. if self._logger:
  15. self._logger(f"[Safe-Redis] {msg}")
  16. else:
  17. print(f"[Safe-Redis] {msg}")
  18. def _safe_execute(self, func, default_return, *args, **kwargs):
  19. """通用安全执行器:遇错自动重试一次,彻底失败返回安全默认值"""
  20. try:
  21. return func(*args, **kwargs)
  22. except Exception as e:
  23. func_name = func.__name__ if hasattr(func, '__name__') else 'operation'
  24. self._log(f"Attempt 1 failed ({func_name}): {e}. Retrying in 0.5s...")
  25. time.sleep(0.5)
  26. try:
  27. return func(*args, **kwargs)
  28. except Exception as e2:
  29. self._log(f"Attempt 2 failed ({func_name}): {e2}. Returning safe fallback: {default_return}.")
  30. return default_return
  31. def exists(self, name) -> bool:
  32. return bool(self._safe_execute(self._client.exists, False, name))
  33. def get(self, name) -> Any:
  34. return self._safe_execute(self._client.get, None, name)
  35. def set(self, name, value, ex=None) -> bool:
  36. return self._safe_execute(self._client.set, False, name, value, ex=ex)
  37. def setex(self, name, time_s, value) -> bool:
  38. return self._safe_execute(self._client.setex, False, name, time_s, value)
  39. def delete(self, *names) -> int:
  40. if not names: return 0
  41. return self._safe_execute(self._client.delete, 0, *names)
  42. def incr(self, name, amount=1) -> int:
  43. # 如果增加失败,默认返回 1,防止触发外部数学运算崩溃(同时保证退避策略能生效)
  44. return self._safe_execute(self._client.incr, 1, name, amount=amount)
  45. def zadd(self, name, mapping) -> int:
  46. if not mapping: return 0
  47. return self._safe_execute(self._client.zadd, 0, name, mapping)
  48. def zrem(self, name, *values) -> int:
  49. if not values: return 0
  50. return self._safe_execute(self._client.zrem, 0, name, *values)
  51. def bulk_zadd(self, name, mapping) -> bool:
  52. """封装 Pipeline,用于一次性安全地写入大量 zset 元素"""
  53. if not mapping:
  54. return True
  55. def _op():
  56. pipe = self._client.pipeline()
  57. for k, v in mapping.items():
  58. pipe.zadd(name, {k: v})
  59. pipe.execute()
  60. return True
  61. return self._safe_execute(_op, False)