You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 23:15:01