MariaDB数据变更时同步数据到SSMS的实现方案咨询
MariaDB 增量变更同步至SQL Server解决方案
以下方案针对无固定更新周期、避免全量加载的需求设计,其中提及的SSMS为SQL Server管理工具,同步目标为SSMS管控的SQL Server数据库。
可行管道设计思路
核心逻辑是跳过全量同步,仅捕获增量变更,两类落地方案适配不同场景:
- 近实时同步方案:开启MariaDB binlog,解析binlog获取每行的增/删/改事件,直接触发同步到SQL Server,延迟在秒级,适合对实时性要求高的场景
- 轻量轮询方案:给需要同步的表加
last_updated时间字段(设置行更新时自动刷新时间戳),小间隔轮询该字段抓取上次同步后的变更数据,适合无法开启binlog的场景,可自行调整轮询间隔平衡资源占用和实时性
Azure Data Factory 实现方案
基于原生CDC的近实时方案
- 前置准备
- 开启MariaDB的binlog,格式设置为ROW模式(必须,才能捕获行级变更)
- 在ADF中分别新建MariaDB和目标SQL Server的链接服务,确保网络连通
- 管道配置
- 新建变更数据捕获(CDC)资源,数据源选MariaDB链接服务,自动检测binlog中的变更事件,目标选SQL Server链接服务,自动映射表结构
- 配置轮转窗口触发器,设置1-5分钟的短间隔(可按需调整),触发CDC任务拉取未处理的变更数据
- 可选配置容错:添加错误日志输出到Azure Blob存储,异常时自动重试2-3次
- 全量初始化处理
- 首次同步时单独跑一次全量复制任务,后续CDC只会同步增量变更,不会产生多余全量加载
备选轻量方案(无法开启binlog时使用):使用上文提到的
last_updated字段,ADF管道配置为:第一步用Lookup活动查询上次同步的最大时间戳,第二步用Copy活动拉取源表中last_updated大于该时间的行写入目标表,配置短间隔触发器即可实现增量同步。
Python 自定义实现代码
基于binlog解析的近实时同步代码,依赖python-mysql-replication解析binlog、pyodbc写入SQL Server:
- 先安装依赖包
pip install python-mysql-replication pyodbc
- 同步代码示例
from pymysqlreplication import BinLogStreamReader from pymysqlreplication.row_event import DeleteRowsEvent, UpdateRowsEvent, WriteRowsEvent import pyodbc # 配置参数,替换为实际信息 MARIA_DB_CONFIG = { "host": "MariaDB服务地址", "port": 3306, "user": "登录账号", "passwd": "登录密码" } SQL_SERVER_CONN_STR = "DRIVER={ODBC Driver 17 for SQL Server};SERVER=SQL Server地址;DATABASE=目标库名;UID=账号;PWD=密码" TARGET_TABLE = "SQL Server目标表名" SOURCE_TABLE = "MariaDB源表名" SERVER_ID = 1 # 自定义唯一ID,不要和MariaDB其他从节点ID冲突 PRIMARY_KEY = "id" # 替换为同步表的实际主键 # 初始化SQL Server连接 sql_conn = pyodbc.connect(SQL_SERVER_CONN_STR) cursor = sql_conn.cursor() # 启动binlog持续监听 stream = BinLogStreamReader( connection_settings=MARIA_DB_CONFIG, server_id=SERVER_ID, only_events=[DeleteRowsEvent, WriteRowsEvent, UpdateRowsEvent], only_tables=[SOURCE_TABLE], blocking=True, # 无新事件时阻塞等待 resume_stream=True # 断线重连后从上次消费位点继续处理 ) for binlogevent in stream: for row in binlogevent.rows: # 处理插入事件 if isinstance(binlogevent, WriteRowsEvent): data = row["values"] cols = ", ".join(data.keys()) placeholders = ", ".join(["?"] * len(data)) sql = f"INSERT INTO {TARGET_TABLE} ({cols}) VALUES ({placeholders})" cursor.execute(sql, list(data.values())) # 处理更新事件 elif isinstance(binlogevent, UpdateRowsEvent): after_data = row["after_values"] pk_val = after_data.pop(PRIMARY_KEY) set_clause = ", ".join([f"{k} = ?" for k in after_data.keys()]) sql = f"UPDATE {TARGET_TABLE} SET {set_clause} WHERE {PRIMARY_KEY} = ?" cursor.execute(sql, list(after_data.values()) + [pk_val]) # 处理删除事件 elif isinstance(binlogevent, DeleteRowsEvent): data = row["values"] pk_val = data[PRIMARY_KEY] sql = f"DELETE FROM {TARGET_TABLE} WHERE {PRIMARY_KEY} = ?" cursor.execute(sql, [pk_val]) sql_conn.commit() stream.close() sql_conn.close()
注意:生产环境使用时建议将binlog消费位点持久化到本地文件或独立数据库,避免进程重启后重复消费数据。
内容的提问来源于stack exchange,提问作者Yash Patil
相关产品推荐
相关产品推荐

