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

如何在Python中感知MariaDB数据库数据变更并触发回调

实现MariaDB数据变更通知与Python缓存刷新

针对你的技术栈(Python 2.7 + MariaDB 10.3 + mysql-connector 2.1.6),以下是几种无需依赖轮询标志表的可行方案:

方案1:基于MariaDB Binlog的变更捕获(CDC)

MariaDB的二进制日志(binlog)会记录所有数据变更操作,你可以通过解析binlog实时监听插入/更新事件,触发缓存刷新。

实现步骤:

  1. 开启MariaDB binlog:
    修改my.cnf配置文件,添加或确认以下配置:

    [mysqld]
    log_bin = /var/log/mysql/mariadb-bin
    binlog_format = ROW  # 行级格式,可捕获具体行变更
    server_id = 1  # 唯一ID,避免与其他复制节点冲突
    

    重启MariaDB服务生效。

  2. 创建具备binlog读取权限的用户:

    GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'binlog_reader'@'localhost' IDENTIFIED BY 'your_password';
    FLUSH PRIVILEGES;
    
  3. 编写Python监听代码:
    对于Python 2.7,可安装兼容版本的mysql-replication库(pip install mysql-replication==0.21.0),编写监听逻辑:

    from pymysqlreplication import BinLogStreamReader
    from pymysqlreplication.row_event import WriteRowsEvent, UpdateRowsEvent
    
    def refresh_cache():
        # 替换为你的缓存重新加载逻辑
        print("Refreshing cache due to data change...")
        # 示例:重新查询数据并更新内存缓存
        # conn = connect(...)
        # cursor = conn.cursor()
        # cursor.execute("SELECT * FROM your_target_table")
        # global cache
        # cache = cursor.fetchall()
    
    def main():
        stream = BinLogStreamReader(
            connection_settings={
                "host": "localhost",
                "port": 3306,
                "user": "binlog_reader",
                "passwd": "your_password"
            },
            server_id=2,  # 需与MariaDB的server_id不同
            blocking=True,  # 持续监听事件
            only_events=[WriteRowsEvent, UpdateRowsEvent]  # 仅监听插入、更新事件
        )
    
        for event in stream:
            # 仅处理目标表的变更
            if event.table == "your_target_table":
                refresh_cache()
    
    if __name__ == "__main__":
        main()
    

    可将该监听逻辑集成到主服务进程中,或通过进程间通信(如本地socket、队列)通知主服务刷新缓存。

方案2:触发器结合外部通知(谨慎使用)

通过MariaDB触发器,在数据变更时调用外部脚本发送通知,触发主服务刷新缓存。此方案需依赖lib_mysqludf_sys用户定义函数(UDF)执行外部命令。

实现步骤:

  1. 安装并注册lib_mysqludf_sys:
    以Debian/Ubuntu为例:

    apt-get install libmysqludf-sys
    

    登录MariaDB注册UDF:

    CREATE FUNCTION sys_exec RETURNS INT SONAME 'libmysqludf_sys.so';
    
  2. 创建数据变更触发器:

    DELIMITER //
    CREATE TRIGGER trigger_after_insert AFTER INSERT ON your_target_table
    FOR EACH ROW
    BEGIN
        SET @cmd = 'python /path/to/notify_refresh.py';
        SELECT sys_exec(@cmd);
    END //
    
    CREATE TRIGGER trigger_after_update AFTER UPDATE ON your_target_table
    FOR EACH ROW
    BEGIN
        SET @cmd = 'python /path/to/notify_refresh.py';
        SELECT sys_exec(@cmd);
    END //
    DELIMITER ;
    
  3. 编写通知脚本notify_refresh.py:
    通过本地socket向主服务发送刷新信号:

    import socket
    
    def send_refresh_signal():
        s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
        try:
            s.connect('/tmp/cache_refresh.sock')
            s.sendall('REFRESH'.encode())
        except Exception as e:
            print("Failed to send refresh signal:", e)
        finally:
            s.close()
    
    if __name__ == "__main__":
        send_refresh_signal()
    

    主服务需监听该socket,收到信号后执行refresh_cache()。

    注意:此方案存在安全风险,UDF可执行任意系统命令,需严格控制数据库权限,仅在可信环境中使用。

方案3:优化版轮询(过渡方案)

若以上方案暂无法实施,可优化你最初的轮询逻辑,减少开销:

  • 直接查询目标表的更新时间,无需额外标志表:
    import time
    from mysql.connector import connect
    
    def get_table_update_time():
        conn = connect(host='localhost', user='your_user', passwd='your_pass', db='your_db')
        cursor = conn.cursor()
        cursor.execute("SELECT UPDATE_TIME FROM INFORMATION_SCHEMA.TABLES WHERE TABLE_SCHEMA = %s AND TABLE_NAME = %s", ('your_db', 'your_target_table'))
        result = cursor.fetchone()
        conn.close()
        return result[0] if result else None
    
    last_update_time = get_table_update_time()
    
    while True:
        current_update_time = get_table_update_time()
        if current_update_time != last_update_time:
            refresh_cache()
            last_update_time = current_update_time
        time.sleep(60)  # 可根据需求调整检查间隔
    
  • 或在目标表中添加last_modified时间戳字段,每次变更时自动更新,轮询时仅检查该字段是否大于上次缓存时间。

内容的提问来源于stack exchange,提问作者Fabiano

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 10:30:52