如何用Python将Redis缓存存入MongoDB并实现实时同步与过期消息处理
Python 操作 Redis 与 MongoDB 的实现方案
1. 将 Redis 缓存数据同步至 MongoDB
首先安装依赖库:
pip install redis pymongo
实现代码
import redis from pymongo import MongoClient import datetime # Redis连接配置 redis_conn = redis.Redis( host='localhost', port=6379, username='admin', password='your-redis-pass', db=0, decode_responses=True ) # MongoDB连接配置 mongo_client = MongoClient('mongodb://localhost:27017/') mongo_db = mongo_client['cache_archive'] mongo_col = mongo_db['redis_cache'] # 同步所有String类型缓存数据 def sync_redis_to_mongo(): # 用scan代替keys避免Redis阻塞(大数据量场景) cursor = 0 while True: cursor, keys = redis_conn.scan(cursor, match='*', count=100) for key in keys: try: value = redis_conn.get(key) doc = { 'redis_key': key, 'value': value, 'sync_time': datetime.datetime.now() } # 存在则更新,不存在则插入 mongo_col.update_one({'redis_key': key}, {'$set': doc}, upsert=True) print(f"同步成功: {key}") except Exception as e: print(f"同步失败 {key}: {str(e)}") if cursor == 0: break # 同步Hash类型数据(按需调整) def sync_redis_hash_to_mongo(): cursor = 0 while True: cursor, keys = redis_conn.scan(cursor, match='user:*', count=100) for key in keys: hash_data = redis_conn.hgetall(key) doc = { 'redis_key': key, 'hash_data': hash_data, 'sync_time': datetime.datetime.now() } mongo_col.update_one({'redis_key': key}, {'$set': doc}, upsert=True) if cursor == 0: break if __name__ == '__main__': sync_redis_to_mongo()
注意事项
- 大数据量场景必须用
scan遍历key,避免keys命令阻塞Redis - 根据实际缓存结构(String/Hash/List等)调整数据读取逻辑
2. 订阅 Redis 过期消息并写入 MongoDB
步骤1:开启 Redis 键空间通知
修改Redis配置文件redis.conf(永久生效):
notify-keyspace-events Ex
或通过命令行临时开启(重启Redis后失效):
redis-cli config set notify-keyspace-events Ex
步骤2:Python 实现订阅与写入
import redis from pymongo import MongoClient import datetime # Redis连接 redis_conn = redis.Redis( host='localhost', port=6379, username='admin', password='your-redis-pass', db=0, decode_responses=True ) # MongoDB连接 mongo_client = MongoClient('mongodb://localhost:27017/') mongo_db = mongo_client['cache_events'] expire_col = mongo_db['expired_keys'] def subscribe_expired_events(): # 订阅对应数据库的过期事件频道(db编号需与Redis连接一致) pubsub = redis_conn.pubsub() pubsub.subscribe('__keyevent@0__:expired') print("开始监听Redis过期事件...") for message in pubsub.listen(): if message['type'] == 'message': expired_key = message['data'] # 过期后key已被删除,无法直接获取数据;需提前在设置key时同步备份 doc = { 'expired_key': expired_key, 'expire_time': datetime.datetime.now() } expire_col.insert_one(doc) print(f"已记录过期key: {expired_key}") if __name__ == '__main__': subscribe_expired_events()
注意事项
- 若需保存过期缓存的内容,建议在设置Redis key时同时将数据写入MongoDB临时集合,过期时再迁移至归档集合
- 高可靠性场景可搭配Redis Stream或消息队列,避免订阅中断丢失事件
3. Redis 实例连接配置
以下是redis-py库的完整连接示例,包含所有核心参数:
import redis # 直接参数配置 redis_conn = redis.Redis( host='your-redis-host', # 如localhost或远程IP port=6379, # 默认端口 username='redis-user', # Redis 6.0+支持用户名,无需可省略 password='your-strong-password', # 无密码可省略 db=0, # 数据库编号,范围0-15,默认0 decode_responses=True, # 自动解码二进制数据为字符串,可选 socket_timeout=5, # 连接超时时间(秒) socket_connect_timeout=5 ) # 或使用URL格式连接 redis_url = 'redis://:your-password@your-redis-host:6379/0' redis_conn = redis.from_url(redis_url, decode_responses=True) # 测试连接 try: redis_conn.ping() print("Redis连接成功") except Exception as e: print(f"Redis连接失败: {str(e)}")
说明
- 启用SSL的Redis实例需添加
ssl=True参数 - 集群环境需使用
redis.RedisCluster类,参数配置略有差异
4. 实时获取插入 Redis 的数据
方案1:键空间通知监听写入事件
开启Redis键空间通知(修改配置为notify-keyspace-events Kx),然后订阅写入事件:
import redis from pymongo import MongoClient import datetime redis_conn = redis.Redis( host='localhost', port=6379, username='admin', password='your-redis-pass', db=0, decode_responses=True ) mongo_client = MongoClient('mongodb://localhost:27017/') mongo_db = mongo_client['real_time_cache'] insert_col = mongo_db['new_entries'] def listen_insert_events(): pubsub = redis_conn.pubsub() # 监听set操作事件,可按需添加hset、lpush等事件 pubsub.subscribe('__keyspace@0__:set') print("开始监听Redis插入事件...") for message in pubsub.listen(): if message['type'] == 'message': inserted_key = message['data'] value = redis_conn.get(inserted_key) doc = { 'key': inserted_key, 'value': value, 'insert_time': datetime.datetime.now() } insert_col.insert_one(doc) print(f"捕获新插入数据: {inserted_key}") if __name__ == '__main__': listen_insert_events()
方案2:Redis Stream 实现可靠实时消费
适合业务可控的场景,支持消息持久化和重复消费:
# 生产者:插入数据到Redis Stream def produce_to_stream(key, value): redis_conn.xadd('cache_stream', {'key': key, 'value': value}) # 消费者:实时读取Stream数据 def consume_stream(): last_id = '0' while True: # 阻塞读取新消息,block=0表示永久阻塞 messages = redis_conn.xread({'cache_stream': last_id}, block=0) for stream, entries in messages: for entry_id, data in entries: doc = { 'stream_entry_id': entry_id, 'key': data['key'], 'value': data['value'], 'consume_time': datetime.datetime.now() } insert_col.insert_one(doc) last_id = entry_id print(f"处理Stream消息: {entry_id}")
注意事项
- 键空间通知适合监听全量写入操作,但可靠性较低(订阅中断会丢失消息)
- Redis Stream适合需要高可靠的业务场景,支持消息回溯和消费确认
内容的提问来源于stack exchange,提问作者Surya RE
相关产品推荐
相关产品推荐

