如何在Python中感知MariaDB数据库数据变更并触发回调
针对你的技术栈(Python 2.7 + MariaDB 10.3 + mysql-connector 2.1.6),以下是几种无需依赖轮询标志表的可行方案:
方案1:基于MariaDB Binlog的变更捕获(CDC)
MariaDB的二进制日志(binlog)会记录所有数据变更操作,你可以通过解析binlog实时监听插入/更新事件,触发缓存刷新。
实现步骤:
开启MariaDB binlog:
修改my.cnf配置文件,添加或确认以下配置:[mysqld] log_bin = /var/log/mysql/mariadb-bin binlog_format = ROW # 行级格式,可捕获具体行变更 server_id = 1 # 唯一ID,避免与其他复制节点冲突重启MariaDB服务生效。
创建具备binlog读取权限的用户:
GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'binlog_reader'@'localhost' IDENTIFIED BY 'your_password'; FLUSH PRIVILEGES;编写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)执行外部命令。
实现步骤:
安装并注册
lib_mysqludf_sys:
以Debian/Ubuntu为例:apt-get install libmysqludf-sys登录MariaDB注册UDF:
CREATE FUNCTION sys_exec RETURNS INT SONAME 'libmysqludf_sys.so';创建数据变更触发器:
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 ;编写通知脚本
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

