数据先接进来:一套可复用的MySQL数据接入与同步方案

发布时间:2026/9/2 19:36:15
数据先接进来:一套可复用的MySQL数据接入与同步方案 不知道你有没有遇到过这样的项目业务方说系统下周就要上线但数据模型还在频繁调整报表需求每天都在变。如果你坚持等所有表结构稳定了再开始做数据同步项目大概率会延期。后来我在多个数据接入项目里逐渐形成一个原则数据先接进来结构后面再调。先把源系统的数据原样同步到目标存储让业务能基于真实数据跑起来再逐步梳理清洗逻辑和建模规范。这套思路看起来简单实际落地时涉及的细节并不少。这篇文章就围绕“数据先接进来”这个思路分享一套可以复用的数据接入方案包含核心设计思路、完整可运行的 Python 代码、常见踩坑点和工程建议。无论是做数据仓库、数据中台还是简单的跨库数据同步都能直接参考。1. 背景与核心概念1.1 什么是“数据先接进来”“数据先接进来”并不是一个严格的技术术语更像是一种工程策略。它的核心思想是在业务模型尚不稳定、下游消费方尚未完全确定的时候先把源系统的数据完整、准时地同步到一个统一位置之后再慢慢处理数据质量问题。这句话可以拆成几个层面理解数据接入先行不依赖下游表结构完全定稿先把源表原样同步过来。保留原始信息接入阶段尽量不做过多的转换和清洗最大程度保留源数据的完整性。为后续建模留空间数据落到统一存储后后续可以随时基于这份原始数据做清洗、宽表加工、指标计算。很多团队在项目初期容易被模型设计拖住。明明业务方只需要先看一下线上真实订单数据但数据团队花了两周讨论数仓分层和字段命名最后连数据都还没接进来。这不是技术问题而是推进节奏的问题。1.2 为什么先接数据而不是先做模型先接数据有几个实际好处第一降低返工成本。数据模型如果一开始设计得很重后续业务规则一变整个链路都要改。而“原样接入”阶段只是把字段原封不动搬过来即使模型调整接入层基本不用动。第二让业务方尽快看到真实数据。业务催数据的时候最怕看到一堆空表和原型数据。先把真实数据接进来哪怕没有复杂加工业务方也能基于原始数据进行初步分析。第三为数据质量摸底留时间。只有先接入数据你才会发现源系统的脏数据、重复数据、字段缺失、类型不一致等问题。这些发现会直接影响后续建模决策。1.3 数据接入的常见应用场景这套思路在以下场景中非常常见数据仓库 ODS 层建设把各业务库的数据同步到数仓作为贴源层。系统迁移从旧系统往新系统迁移数据先全量后增量。多库汇聚把多个分库分表实例的数据汇聚到一张总表。异地容灾备份把核心业务数据定时同步到灾备库。实时数仓的离线兜底实时链路不稳定时用离线同步作为数据补偿。1.4 自研同步工具还是用现成框架说到数据接入很多人会想到 DataX、Kettle、Canal、Flink CDC 等现成工具。那为什么还要自己写现成框架适合大规模、复杂场景但它们也有学习成本和部署成本。对于中小团队、单体应用、跨库表数量不多的情况一个轻量的 Python 脚本反而更直接也更容易维护。它不需要额外部署服务不需要学习一套新 DSL出现问题也可以直接改代码排查。本文的案例就是一个可运行的轻量同步脚本重点展示数据接入的核心工程设计而不是依赖某个重量级框架。理解这些原理后你再去用 DataX 或 Flink CDC也会更快上手。2. 环境准备与版本说明2.1 运行环境本文示例以 Python 为基础编写你本机需要安装 Python 3.10 及以上版本。具体环境如下操作系统Windows / Linux / macOS 均可Python3.10 或更高MySQL5.7 或 8.0PyMySQL用于 Python 连接 MySQLPyYAML用于读取配置文件安装依赖pip install pymysql pyyaml2.2 示例数据库为了演示完整流程我准备了两套 MySQL 数据库shop_db业务源库存放线上业务表。dw_db目标库模拟数据仓库或数据汇聚层。如果你的环境里只有一套 MySQL可以建两个库来模拟逻辑是一样的。版本方面不需要过分纠结MySQL 5.7 和 8.0 在本文涉及的语法上没有明显差异。2.3 项目目录结构整个同步工具的项目结构如下data_ingest/ ├── config.yaml # 同步配置文件 ├── db.py # 数据库连接工具 ├── sync.py # 同步核心逻辑 ├── main.py # 程序入口 ├── sql/ │ └── init.sql # 源表和目标表初始化脚本 └── README.md # 项目说明可选后面会逐个文件讲解你可以直接复制代码按需修改配置。3. 数据接入的核心设计思路动手写代码之前先梳理几个关键设计点。这些点直接决定了同步工具是否可靠。3.1 全量同步与增量同步数据接入首先要区分两个阶段全量同步和增量同步。全量同步把源表当前的所有数据一次性同步到目标表。通常在首次接入时执行也适用于数据量小、且目标表需要重建的场景。增量同步在全量同步之后按照时间、序号等标记只同步发生变化的数据。适合数据量持续增长的业务表。全量同步有一个明显缺点如果每天跑全量数据量大了以后不仅耗时还会对源库造成较大压力。因此更合理的做法是首次全量之后按增量频率执行。增量同步依赖一个“增量字段”。常见的增量字段有两种增量字段类型示例特点时间字段updated_at、create_time直观但依赖业务系统正确维护更新时间自增主键id只能识别新增无法捕获修改版本号version需要业务系统配合本文示例采用updated_at作为增量字段这也是大多数业务系统最容易提供的字段。3.2 幂等写入数据同步过程中很容易出现重复写入的情况。比如增量任务因为网络超时被重试。上一批数据已经写入但程序在提交状态之前崩溃了。两条不同时间点的数据修改后更新到同一主键。如果目标表只是简单INSERT就会出现主键冲突或重复数据。解决办法是让写入操作具备幂等性也就是说同一行数据无论写入几次最终状态都是一致的。在 MySQL 中最常用的方式就是INSERT ... ON DUPLICATE KEY UPDATE。它检测到唯一键冲突时不会报错而是执行更新操作。这样就能保证同一主键的数据只会保留最新状态。3.3 同步状态记录增量同步最怕一个问题任务中断后从哪里继续解决思路是增加一张同步状态表专门记录每个表“上一次同步到哪个位置”。每次增量任务开始前读取这个状态作为起点任务结束后更新这个状态。有了状态记录同步任务就具备了“断点续传”能力。即使某次任务因为网络故障中断重新运行时也能从上一次成功的位置继续而不是从头再来。3.4 分批拉取如果源表数据量很大一次性SELECT * FROM table会带来两个问题大量数据同时进入内存容易导致 OOM。对源库产生较长时间的查询压力影响线上业务。因此同步逻辑必须分批拉取数据。分批方式有两种常见思路基于 LIMIT 分页SELECT * FROM t_order ORDER BY id LIMIT 0, 5000; SELECT * FROM t_order ORDER BY id LIMIT 5000, 5000;这种方式实现简单但遇到深分页时比如第 100 万条以后MySQL 需要扫描并丢弃大量行性能会明显下降。基于主键范围分批SELECT * FROM t_order WHERE id 1 AND id 5001; SELECT * FROM t_order WHERE id 5001 AND id 10001;这种方式使用主键索引定位范围不会因为偏移量增大而变慢是更推荐的做法。本文代码会采用这种方式。3.5 字段白名单同步工具往往会拼接表名和字段名到 SQL 中。如果不加限制一旦配置文件被恶意修改就可能产生 SQL 注入风险。最佳实践是在代码中维护一份表名白名单、字段白名单只有通过校验的对象才允许拼入 SQL。这样既保证了灵活性也守住了安全底线。4. 完整实战案例接下来进入正题。我以“把shop_db的订单表t_order同步到dw_db的ods_order表”为例完整演示整个数据接入过程。4.1 初始化表结构先执行sql/init.sql创建源表和目标表。-- 文件路径data_ingest/sql/init.sql -- 源库表结构示例 CREATE DATABASE IF NOT EXISTS shop_db DEFAULT CHARSET utf8mb4; USE shop_db; DROP TABLE IF EXISTS t_order; CREATE TABLE t_order ( id BIGINT NOT NULL AUTO_INCREMENT, order_no VARCHAR(64) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(12, 2) NOT NULL, status TINYINT NOT NULL DEFAULT 0, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, PRIMARY KEY (id), KEY idx_updated_at (updated_at) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; -- 目标库表结构 CREATE DATABASE IF NOT EXISTS dw_db DEFAULT CHARSET utf8mb4; USE dw_db; DROP TABLE IF EXISTS ods_order; CREATE TABLE ods_order ( id BIGINT NOT NULL, order_no VARCHAR(64) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(12, 2) NOT NULL, status TINYINT NOT NULL DEFAULT 0, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, sync_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4; -- 同步状态表 DROP TABLE IF EXISTS sync_state; CREATE TABLE sync_state ( table_name VARCHAR(128) NOT NULL, sync_time DATETIME NOT NULL, PRIMARY KEY (table_name) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;这里有几个设计点需要说明目标表比源表多了一个sync_time字段用来记录这行数据是在哪个时间点被同步进来的。这属于接入层的“技术痕迹”不影响业务数据本身。sync_state表存的是每个表“上次同步到的时间点”不是同步任务本身的历史记录。它是增量同步的游标。初始化完成后给源表插入几条模拟订单数据USE shop_db; INSERT INTO t_order (order_no, user_id, amount, status, created_at, updated_at) VALUES (ORD20250101001, 1001, 199.00, 1, NOW(), NOW()), (ORD20250101002, 1002, 299.00, 1, NOW(), NOW()), (ORD20250101003, 1003, 99.50, 0, NOW(), NOW());4.2 编写配置文件创建config.yaml配置源库、目标库和同步参数。# 文件路径data_ingest/config.yaml source: host: 127.0.0.1 port: 3306 user: root password: your_password database: shop_db target: host: 127.0.0.1 port: 3306 user: root password: your_password database: dw_db sync: table: t_order target_table: ods_order incremental_field: updated_at batch_size: 5000配置项说明source源库连接信息。target目标库连接信息。sync.table要同步的源表名。sync.target_table目标表名。sync.incremental_field增量字段名这里使用updated_at。sync.batch_size每批拉取的行数具体值需要根据表数据量和字段宽度调整。配置文件中不要硬编码生产环境的账号密码。实际工程中应该通过环境变量或密钥管理服务注入这个后面在最佳实践部分展开。4.3 编写数据库连接工具创建db.py统一管理数据库连接。# 文件路径data_ingest/db.py import pymysql def get_connection(db_config: dict): 根据数据库配置创建 MySQL 连接。 这里统一使用 DictCursor查询结果会以字典形式返回 后续代码读取列名时会更直观。 return pymysql.connect( hostdb_config[host], portdb_config[port], userdb_config[user], passworddb_config[password], databasedb_config[database], charsetutf8mb4, cursorclasspymysql.cursors.DictCursor, )charsetutf8mb4这一项不能省。如果源数据里包含表情符号或其他扩展字符不指定这个字符集很容易出现中文乱码或写入失败。4.4 编写同步核心逻辑创建sync.py这是整个工具的骨干包含全量同步、增量同步、幂等写入和状态记录。# 文件路径data_ingest/sync.py from datetime import datetime from db import get_connection # 表名白名单只允许同步明确授权的表 ALLOWED_TABLES {t_order} # 字段白名单增量字段必须是这里面的字段 ALLOWED_FIELDS {id, order_no, user_id, amount, status, created_at, updated_at} def check_table_name(table_name: str): if table_name not in ALLOWED_TABLES: raise ValueError(f表 {table_name} 不在白名单中禁止同步) def check_field_name(field_name: str): if field_name not in ALLOWED_FIELDS: raise ValueError(f字段 {field_name} 不在白名单中禁止使用) def get_last_sync_time(target_conn, target_table: str): 读取上次同步更新时间。 如果 sync_state 表中还没有记录说明这个目标表从未同步过 此时返回 None调用方会切换为全量同步。 check_table_name(target_table) with target_conn.cursor() as cursor: cursor.execute( SELECT sync_time FROM sync_state WHERE table_name %s, (target_table,), ) row cursor.fetchone() if row: return row[sync_time] return None def update_sync_time(target_conn, target_table: str, sync_time): 更新同步状态表中的同步时间。 使用 ON DUPLICATE KEY UPDATE避免重复插入状态记录。 check_table_name(target_table) with target_conn.cursor() as cursor: cursor.execute( INSERT INTO sync_state (table_name, sync_time) VALUES (%s, %s) ON DUPLICATE KEY UPDATE sync_time VALUES(sync_time) , (target_table, sync_time), ) target_conn.commit() def batch_insert(target_conn, target_table: str, rows: list, sync_timeNone): 将一批数据写入目标表。 使用 INSERT ... ON DUPLICATE KEY UPDATE 保证幂等 如果主键已存在则更新业务字段而不是报错。 if not rows: return sync_time sync_time or datetime.now() # 给每行补充同步时间字段 enriched_rows [] for row in rows: new_row {**row, sync_time: sync_time} enriched_rows.append(new_row) columns list(enriched_rows[0].keys()) # 主键 id 不参与更新sync_time 由本脚本维护也不从源数据更新 update_cols [col for col in columns if col not in (id, sync_time)] col_sql , .join(columns) placeholders , .join([%s] * len(columns)) update_sql , .join([f{col} VALUES({col}) for col in update_cols]) sql ( fINSERT INTO {target_table} ({col_sql}) fVALUES ({placeholders}) fON DUPLICATE KEY UPDATE {update_sql} ) values [[row[col] for col in columns] for row in enriched_rows] with target_conn.cursor() as cursor: cursor.executemany(sql, values) target_conn.commit() def full_sync(source_conn, target_conn, table_name, target_table, batch_size): 全量同步。 流程 1. 清空目标表保证目标表与源表一致。 2. 按主键 id 范围分批拉取源表数据。 3. 每批数据写入目标表。 4. 更新同步状态。 check_table_name(table_name) check_table_name(target_table) # 清空目标表 with target_conn.cursor() as cursor: cursor.execute(fTRUNCATE TABLE {target_table}) target_conn.commit() sync_time datetime.now() with source_conn.cursor() as cursor: # 获取主键范围 cursor.execute(fSELECT MIN(id) AS min_id, MAX(id) AS max_id FROM {table_name}) row cursor.fetchone() min_id row[min_id] max_id row[max_id] if min_id is None: update_sync_time(target_conn, target_table, sync_time) print(f[全量同步] 表 {table_name} 没有数据跳过) return total_count 0 current_id min_id # 按主键范围分批拉取避免深分页性能问题 while current_id max_id: end_id current_id batch_size cursor.execute( fSELECT * FROM {table_name} WHERE id %s AND id %s, (current_id, end_id), ) rows cursor.fetchall() if rows: batch_insert(target_conn, target_table, rows, sync_time) total_count len(rows) print(f[全量同步] 已同步 {len(rows)} 条id 范围{current_id} ~ {end_id}) current_id end_id update_sync_time(target_conn, target_table, sync_time) print(f[全量同步] 完成共同步 {total_count} 条数据) def incremental_sync(source_conn, target_conn, table_name, target_table, incremental_field, batch_size): 增量同步。 流程 1. 读取上次同步时间。 2. 如果没有同步记录自动切换为全量同步。 3. 按增量字段排序分批拉取上次同步时间之后的数据。 4. 每批数据写入目标表并推进游标时间。 5. 全部完成后更新同步状态。 check_table_name(table_name) check_table_name(target_table) check_field_name(incremental_field) last_sync_time get_last_sync_time(target_conn, target_table) if last_sync_time is None: print([增量同步] 未找到上次同步时间自动切换为全量同步) full_sync(source_conn, target_conn, table_name, target_table, batch_size) return sync_time datetime.now() total_count 0 current_time last_sync_time with source_conn.cursor() as cursor: while True: cursor.execute( f SELECT * FROM {table_name} WHERE {incremental_field} %s ORDER BY {incremental_field} ASC LIMIT %s , (current_time, batch_size), ) rows cursor.fetchall() if not rows: break batch_insert(target_conn, target_table, rows, sync_time) total_count len(rows) # 用最后一条数据的增量字段值作为新的游标 current_time rows[-1][incremental_field] print(f[增量同步] 已同步 {len(rows)} 条游标时间{current_time}) update_sync_time(target_conn, target_table, sync_time) print(f[增量同步] 完成共处理 {total_count} 条数据)这段代码有几个细节值得展开说明。为什么全量同步要先 TRUNCATE 目标表TRUNCATE 比 DELETE 速度更快而且会重置自增 ID。全量同步的目标就是让目标表和源表完全一致先清空再写入是最直接的方式。不过要注意TRUNCATE 不能回滚操作前必须确认目标表确实可以被清空。为什么批量插入用executemanyexecutemany会批量执行 INSERT减少客户端和数据库之间的交互次数。如果一条一条执行几万条数据会产生大量网络往返性能差距非常明显。为什么增量同步中游标时间用最后一条数据的时间增量任务拉取数据后游标应该推进到“已经处理的位置”。用最后一条数据的updated_at作为新游标可以避免每批重复扫描之前的数据。但这里有一个边界情况如果同一秒内有多条数据的updated_at相同且刚好分布在两批之间current_time推进后下一批使用 current_time会漏掉同一时间戳的剩余数据。解决方法是改用复合游标(updated_at, id)或者在业务上保证updated_at的精度足够高。这个在常见问题部分会再讲。4.5 编写程序入口创建main.py统一调度全量和增量同步。# 文件路径data_ingest/main.py import argparse import yaml from db import get_connection from sync import full_sync, incremental_sync def load_config(config_path: str) - dict: with open(config_path, r, encodingutf-8) as f: return yaml.safe_load(f) def main(): parser argparse.ArgumentParser(description数据接入同步工具) parser.add_argument(--config, defaultconfig.yaml, help配置文件路径) parser.add_argument( --mode, choices[auto, full, incremental], defaultauto, help同步模式auto 自动判断full 强制全量incremental 强制增量, ) args parser.parse_args() config load_config(args.config) source_conn get_connection(config[source]) target_conn get_connection(config[target]) table_name config[sync][table] target_table config[sync][target_table] incremental_field config[sync][incremental_field] batch_size config[sync][batch_size] try: if args.mode full: full_sync(source_conn, target_conn, table_name, target_table, batch_size) elif args.mode incremental: incremental_sync( source_conn, target_conn, table_name, target_table, incremental_field, batch_size, ) else: # auto 模式 # 增量同步内部会判断是否存在上次同步时间 # 如果没有则自动切换为全量同步。 incremental_sync( source_conn, target_conn, table_name, target_table, incremental_field, batch_size, ) finally: source_conn.close() target_conn.close() if __name__ __main__: main()--mode参数的设计参考了很多定时任务的用法第一次接入时执行--mode full。后续定时任务执行--mode incremental。如果不想关心当前状态直接执行--mode auto脚本会自动判断。4.6 运行与验证先执行全量同步cd data_ingest python main.py --config config.yaml --mode full预期输出[全量同步] 已同步 3 条id 范围1 ~ 5000 [全量同步] 完成共同步 3 条数据登录目标库验证数据USE dw_db; SELECT * FROM ods_order; SELECT * FROM sync_state;此时ods_order中应该有 3 条数据sync_state中记录了一个同步时间。接着模拟业务新增和修改数据USE shop_db; -- 新增一条订单 INSERT INTO t_order (order_no, user_id, amount, status, created_at, updated_at) VALUES (ORD20250101004, 1004, 599.00, 1, NOW(), NOW()); -- 修改一条已有订单 UPDATE t_order SET amount 259.00, updated_at NOW() WHERE id 1;然后执行增量同步python main.py --config config.yaml --mode incremental预期输出大致如下[增量同步] 已同步 2 条游标时间2025-01-01 10:30:00 [增量同步] 完成共处理 2 条数据再次查看目标表USE dw_db; SELECT * FROM ods_order ORDER BY id;可以看到新增的订单出现了id1 的订单金额也从 199.00 变成了 259.00。这就是增量同步的效果。再执行一次增量同步python main.py --config config.yaml --mode incremental预期输出[增量同步] 完成共处理 0 条数据因为从上次游标时间到现在没有新增或修改数据所以处理 0 条是正确结果。5. 常见问题与排查思路数据接入工具本身逻辑不复杂但在真实环境中会遇到各种边界情况。下面整理了一些高频问题。问题现象常见原因解决思路目标表写入中文乱码Python 连接未指定charsetutf8mb4检查db.py连接配置统一使用utf8mb4增量同步漏数据updated_at时间精度不够同一时刻多条数据分布在不同批次改用(updated_at, id)复合游标或提高时间精度增量同步重复数据游标时间推进方式不正确批次边界重叠使用上一批最后一条记录的字段值作为新游标条件使用同步任务很慢全量拉取使用LIMIT深分页或批次大小不合理改用主键范围分批适当调大batch_size主键冲突导致任务失败目标表缺少幂等写入逻辑使用INSERT ... ON DUPLICATE KEY UPDATE全量同步误清空目标表手工执行脚本时指定了错误的表配置白名单机制禁止同步未授权的表任务中断后重复写入大量数据状态表未更新确认每次任务结束都调用update_sync_time下面单独展开两个最常踩的坑。坑一增量游标时间精度不足导致漏数据假设updated_at精确到秒某一秒内同时有 3 条订单被修改。第一批拉取了一条游标被推进到这个秒值。第二批查询条件变成updated_at 这个秒值剩下的 2 条就被漏掉了。解决办法有两种查询时额外带一个主键条件使用复合游标WHERE (updated_at %s) OR (updated_at %s AND id %s) ORDER BY updated_at, id如果业务量不大可以在下一批查询时使用并用主键去重。但总体上复合游标更可靠。**坑二全量同步过程中