尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

数据同步中间件实战:MySQL到Kafka及多源异构场景的选型与避坑

数据同步中间件实战:MySQL到Kafka及多源异构场景的选型与避坑 简介DBSyncerdbs是一款面向开发与运维人员的数据同步中间件专注于解决多源异构数据库间的实时数据流转与同步难题覆盖MySQL、Oracle、SqlServer、PostgreSQL、Elasticsearch、Kafka、File、SQL等多种主流同步场景并提供插件扩展能力以支持自定义转换业务。压缩包共737个文件大小约2.07MB以472个Java源码文件为核心辅以81个png图片、59个css样式、42个html页面、36个js脚本及xml、sql等配置资源整体属于轻量级中间件工程包适合研究源码结构或快速部署学习。资源内置全量与增量数据统计图展示、应用性能预警等监控能力配合插件机制可实现定制化同步逻辑对需要保障数据同步稳定性的项目有直接借鉴价值。目前已有747人学习下载适合正在调研或落地数据同步方案的开发者参考。1. 数据同步中间件不是ETL它是把多源异构数据搬到目标端的“传送带”“数据同步中间件”听起来是个高深组件实际解决的是最老实的搬运问题MySQL的一张订单表要实时进KafkaOracle的维表要按天整库搬到SqlServerPostgreSQL的变更要落成文件归档。真去手写脚本你得自己处理分页、幂等、断点重跑还要面对binlog这类黑匣子。开源数据同步中间件把源端读取和目标端写入抽象成Source与Sink插件一份任务配置就能覆盖MySQL、Oracle、SqlServer、PostgreSQL、File、Kafka、SQL等同步场景自带checkpoint断点续传。适合数据工程师、运维以及被“数据搬家”折磨过的后端开发但在动手前得先看清它的边界到底在哪。2. 选型先看同步场景MySQL、Oracle、SqlServer、PostgreSQL、File、Kafka、SQL各自的门槛在哪里2.1 全量同步与增量同步“同步”不是一个动作是两套完全不同的机制很多新人接到需求第一反应是“定时把表查出来插到目标库”。这个理解在全量同步里没错但增量同步完全是另一回事。全量同步做的是“当前快照”流程通常是清空目标表或重建表、按分区键分批SELECT、再写入增量同步做的是“变更流”中间件要通过数据库日志捕获变更——MySQL是binlogOracle是LogMinerSqlServer是CDCPostgreSQL是逻辑复制——再按事务顺序投递到目标端。两者在权限要求、数据结构、错误处理上差异很大。做这个标题下的场景必须接受“一个任务先全量、再自动切增量”的模型。开源数据同步中间件在这个问题上的取舍也明显DataX这类工具偏全量批式跑对齐没问题但拿不到delete事件Flink CDC偏实时流秒级感知变更但文件、多库统一收口不是它的最强项Apache SeaTunnel这类中间件能把initial全量加streaming增量放进同一个Pipeline也是我在这类混合场景里用得最多的落地方式。选型逻辑其实简单要跑批对齐选全量强的要秒级感知选CDC强的两者都要就选带initial模式的CDC插件。2.2 JDBC直连、CDC订阅、文件解析七类数据源的三种接入方式标题里的七类场景接入方式可以归成三种。第一种是JDBC直连MySQL、Oracle、SqlServer、PostgreSQL以及SQL查询都能用jdbc url加select语句读取好处是通用、只要读权限就能跑坏处是增量水位要自己维护对源库压力也大第二种是CDC订阅通过数据库日志解析变更能拿到update前的旧值和delete事件对源库侵入小代价是要创建专用账号、开启归档或CDC参数还要处理DDL第三种是文件与消息File走目录轮询Kafka走消费组本质都是流式读取外部事件关键在于位点管理和幂等。接入方式代表源延迟权限与成本适合场景JDBC直连MySQL、Oracle、SqlServer、PostgreSQL、SQL查询分钟级低读权限即可小时级批同步、手动补数CDC订阅MySQL binlog、Oracle LogMiner、SqlServer CDC、PG逻辑复制秒级需开启归档或CDC、专用账号实时增量、交易类数据文件与消息LocalFile、Kafka秒级到准实时低日志汇聚、事件流、数据湖这三种方式不是单选题。同一个任务里重要表走CDC增量、历史表走JDBC全量是很常见的混用。网上经常讨论的“使用flink实现mysql同步到clickhouse”本质也只是实时流里的一种路径而标题覆盖的这批场景更常见的是用SeaTunnel这类工具把JDBC、File、Kafka统一收口。另外提醒一句很多同步失败源于查询里的隐式转换比如把VARCHAR当数字比较导致索引失效这在SqlServer和Oracle的增量筛选里都特别常见。2.3 从源到目标一张表看清七类同步场景的可行路径源目标可行模式关键点MySQLKafkabinlog CDCserver-id唯一、binlog_formatROWOracleSqlServerJDBC全量或LogMiner增量ojdbc驱动版本、分区键SqlServerMySQLCDC或JDBC表级CDC、Agent作业PostgreSQLKafka逻辑复制slot.name唯一防止WAL膨胀FilePostgreSQL目录轮询解析格式、文件完成后renameKafkaMySQL流式消费幂等、主键冲突策略SQL查询结果任意目标JDBC source加query查询要带分区条件这张表的结论是管道本身大同小异真正难的是两端接入差异。很多读者可能刚用rpm装好MySQL 5.7.44或者在Windows上折腾SqlServer 2016安装失败到了同步这一步又会撞上驱动和权限两堵墙。下面第三章先把最小闭环跑通第四章再逐个场景拆。3. 从部署到跑通首条MySQL到Kafka的同步任务最小可复现配置3.1 部署前准备二进制包、JDK与驱动目录不管选哪款中间件第一步都是下载对应版本的二进制发行包而不是源码编译。编译一次可能要消耗一晚上且依赖版本冲突是出了名的坑。解压后确认目录里有bin、config、plugins这些结构然后把用到的数据库JDBC驱动手动放进加载目录有的发行版放在lib有的放在connectors下的子目录认准你那个版本实际扫描class的路径驱动不加载后面跑起来全是ClassNotFoundException。# 解压并进入安装目录 tar -zxvf apache-seatunnel-2.3.x-bin.tar.gz cd apache-seatunnel-2.3.x-bin # 把驱动jar放进lib目录路径按你的发行版实际情况调整 cp ojdbc8.jar lib/ cp mssql-jdbc-*.jar lib/这段命令背后的逻辑数据同步中间件对数据库的访问全部走JDBC但开源包受许可证限制不会内置Oracle和SqlServer驱动必须手动补充。驱动版本还要跟JDK匹配JDK8用ojdbc8JDK17用ojdbc11SqlServer的mssql-jdbc也要选对应的jre版本混用的话会在建立连接时抛UnsupportedClassVersionError。# 后台启动中间件日志输出到logs目录下 bin/seatunnel-cluster.sh -d # 用curl确认集群起来了端口以启动日志为准 curl -s http://localhost:5801/hazelcast/rest/cluster启动成功只是第一步。我一般会在跑任务前看一眼日志目录有没有生成确认进程真的在监听端口而不是“命令没报错就当成功”。3.2 写第一个PipelineMySQL到Kafka的HOCON配置逐段拆解中间件通用的任务描述方式是“一个Pipeline四段式”env定义运行环境source定义从哪读transform定义怎么改sink定义写到哪。下面这份配置以当前主流版本的写法为例把MySQL一张订单表全量加增量同步到Kafkaenv { parallelism 2 job.mode STREAMING checkpoint.interval 10000 } source { MySQL-CDC { result_table_name orders host 192.168.1.20 port 3306 username cdc_user password Cdc2024 server-id 5401-5410 base-url jdbc:mysql://192.168.1.20:3306/shop?useSSLfalseserverTimezoneAsia/Shanghai table { database shop table-names [shop.t_order] } startup.mode initial } } transform { } sink { Kafka { bootstrap.servers 192.168.1.30:9092 topic dwd_t_order format json semantics exactly-once } }这一段配置的要点是startup.mode设成initial任务启动后会先把表当前数据全量发一遍然后自动切到binlog增量所以不需要你手动分成两个任务server-id给的是一个范围而不是单值因为CDC插件会模拟多个从库连接去并行拉binlog范围必须全网唯一否则会跟其他同步任务互相踢下线base-url里的serverTimezoneAsia/Shanghai是必须的不写的话时间字段会按JVM默认时区解析后面基本都是8小时偏差。Kafka sink里的semanticsexactly-once依赖Kafka事务broker侧要开启事务支持否则任务会在commit阶段抛TimeoutException。如果你的Kafka集群没配事务先把语义降成at-least-once在下游按主键去重比硬开exactly-once更省心。3.3 提交任务与验证除了日志还要看这三个信号配置写好后先用本地模式试跑资源不隔离适合验证语法和连通性bin/seatunnel.sh --config conf/mysql_to_kafka.conf -m local生产环境建议把任务提交到集群模式由中间件统一调度这样能拿到失败重启、checkpoint回溯这些能力。跑起来后不要只看进程在不在我一般确认三个信号第一个信号是日志里出现“Checkpoint completed”或类似的提交记录说明源端读取和sink写入的状态已经被可靠保存第二个信号是消费Kafka目标topic看消息有没有进来bin/kafka-console-consumer.sh \ --bootstrap-server 192.168.1.30:9092 \ --topic dwd_t_order \ --from-beginning \ --max-messages 5第三个信号最直观回到MySQL源库update一条记录1到2秒内Kafka应该出现新的变更消息。如果全量数据到了、增量不到问题基本不在中间件而在binlog配置或账号权限这属于后面避坑章的内容。4. 把更多同步场景串起来Oracle、SqlServer、PostgreSQL、File、Kafka、SQL的配置要点4.1 Oracle与SqlServer作为源JDBC分区与CDC开启的前置条件Oracle接入最常见的做法是先用JDBC全量跑通确认链路没问题再决定要不要上LogMiner增量。下面这份配置用Jdbc源把Oracle一张大表按分区键拆成多个分片并行抽取适合维表、流水表这类需要定期对齐的场景source { Jdbc { url jdbc:oracle:thin://192.168.1.20:1521/ORCLPDB1 driver oracle.jdbc.driver.OracleDriver user sync_user password Sync123 query SELECT ID, ORDER_NO, GMT_CREATED FROM T_ORDER WHERE GMT_CREATED :condition partition_column ID partition_num 4 } }这里最关键的是partition_column和partition_num。中间件会按分区键把查询拆成多个区间并行执行分区键必须用数字型主键或时间戳不要用字符串否则每个分片都要做全表扫描Oracle短时间把CPU打满。另外Oracle查询条件里避免把VARCHAR字段跟数字做隐式比较这会直接废掉索引跟SqlServer里“字符串转数字”导致全表扫描是同一个坑。SqlServer要接增量前提是表级CDC已经开启。常见做法是先对库执行sys.sp_cdc_enable_db再对目标表执行sys.sp_cdc_enable_table然后确认Agent作业在运行。很多人卡在“CDC开了但收不到变更”原因基本都是表没加入capture instance或者Agent作业被禁用。配置层面对应的是SqlServer-CDC源重点给足权限和表名source { SqlServer-CDC { host 192.168.1.18 port 1433 username cdc_user password Cdc2024 database erp table-names [erp.dbo.t_order] } }SqlServer 2016、2017、2019、2022开CDC的机制差异不大但要特别注意版本授权Express版不带Agent作业功能开完CDC也不会有人帮你抓变更这是测试环境最常翻车的地方。4.2 PostgreSQL与File作为源复制槽与目录轮询的最小写法PostgreSQL实时同步走逻辑复制不是查表的增量字段。核心是复制槽配置如下source { PostgreSQL-CDC { host 192.168.1.15 port 5432 username replica_user password Pg2024 database appdb slot.name seatunnel_slot decoding.plugin.name pgoutput table-names [appdb.t_user] } }slot.name在整套PostgreSQL实例里必须唯一因为复制槽是实例级资源。比配置更需要注意的是运维如果目标端停了或者这个任务被删了复制槽不会自动清理WAL日志会一直堆积把磁盘写爆。我见过不止一次PG磁盘告警查到最后都是无人认领的复制槽。切换任务前先handshake旧任务停掉、确认slot的restart_lsn不再往前推进、再启动新任务。File作为源时配置比数据库简单但解析格式要提前想清楚source { LocalFile { path /data/sync_files file_format.type csv schema [ { field id, type bigint } { field order_no, type string } { field amount, type decimal } ] } }目录轮询的常见坑是“文件半写状态”。生产环境里上游写入文件不是原子的中间件轮询时可能读到半个文件。更保险的做法是上游先把文件写成tmp写完再rename成正式后缀中间件只认正式后缀消费完的文件怎么处理也要定好策略是先rename到done目录再删还是中间件直接标记消费位点两种方案都要保证“重复跑不产生重复数据”。4.3 从Kafka和SQL查询到目标库反向链路与“SQL即数据源”Kafka到MySQL是实时数仓里最常见的反向链路消费消息、解析JSON、批量写入目标表。配置里sink端用Jdbc自动生成SQLsource { Kafka { bootstrap.servers 192.168.1.30:9092 topic ods_t_order format json } } sink { Jdbc { url jdbc:mysql://192.168.1.20:3306/dw?useSSLfalseserverTimezoneAsia/Shanghai driver com.mysql.cj.jdbc.Driver user dw_user password Dw123 generate_sink_sql true database dw table t_order_dwd batch_size 1000 } }generate_sink_sqltrue的意思是让中间件根据消息里的字段自动生成insert省去手写SQL的琐碎工作。batch_size不建议一味调大我踩过的现实是批量从1000调到5000写入确实快了但源端一有大事务目标端频繁死锁最后又调回2000配合主键冲突策略做upsert才稳定下来。标题里的“SQL同步场景”本质是“把一条SQL查询结果当作同步源”。在Jdbc source里写query中间件按查询结果往目标端灌数这在做宽表加工、数据字典同步时非常实用。唯一要注意的是查询必须能分片也就是query里要带范围条件否则中间件只用单连接全量捞大表会把源库连接池拖垮。5. 数据同步中间件的避坑笔记驱动、时区、server-id与checkpoint的5条血泪经验5.1 Kafka消息延迟高先分清是源头慢还是目标端慢现象任务跑了两小时Kafka目标topic积压几百万条消费组lag只增不减日志里却没有一条报错。原因source端MySQL-CDC拉取正常但sink端Kafka在单条发送且任务并行度只有1。很多中间件的Kafka sink默认不会帮你打开批量发送需要显式设置linger.ms和batch_size这类参数。解决先把并行度提到4再给Kafka sink打开批量发送linger_ms设到50到100毫秒、batch_size设到几百KB的量级观察延迟曲线往下走再加并行度。教训是不要一上来把并行度调到16批量参数没跟上之前并行度越高连接数爆炸越快延迟反而更难看。5.2 SqlServer开启CDC后仍收不到变更现象任务能正常连上全量数据也同步完了但业务表update之后目标库没有任何新数据。原因表没有加入capture instance或者SqlServer的Agent作业没有运行。很多测试库用的是Express版根本不带Agent功能CDC开了等于没开。解决按顺序排查三件事确认对库执行过sys.sp_cdc_enable_db对表执行过sys.sp_cdc_enable_table确认捕获作业在运行任务计划没有被禁用确认同步账号有VIEW ANY DEFINITION权限。这条排查询走完绝大多数“收不到变更”的问题都能解决跟中间件本身无关。5.3 Oracle连接池耗尽同步任务“假活”现象日志反复抛ORA-12519任务重启后好一阵又挂进程都活着但吞吐为零。原因Jdbc source的partition_num设得很大每个分片占用一个连接Oracle的PROCESSES上限被瞬间打满。这个问题容易跟“监听服务无法启动”混淆其实监听没挂是连接数溢出。解决把partition_num降到数据库连接余量能承受的水平或者同步加大数据库processes参数中间件侧如果能配连接池优先复用连接而不是每分片新建。Oracle这类重型库启动任务前先查当前会话数比看日志更直接。5.4 时间字段差8小时时区不是玄学是参数现象源库查询出来时间正常写入目标库后整体慢了8小时有的带TZ字段正常有的不带TZ字段就偏。原因MySQL的jdbc url没写serverTimezone中间件进程默认按UTC解析Oracle的TIMESTAMP WITH TIME ZONE和PostgreSQL的timestamptz语义又各不同统一转成目标库类型时容易错位。解决所有JDBC url统一加上serverTimezoneAsia/Shanghai和useSSLfalseOracle连接串里设置TIME_ZONE08:00。上线之前抽一条记录做“源库值、中间件日志值、目标库值”三方对比能省掉后面一整轮对账排查。5.5 binlog拉取卡死或checkpoint失败多数是server-id占位或大事务现象任务从某个位点后不再前进日志报Disconnect channelcheckpoint反复恢复失败消费位点原地不动。原因两个同步任务配置了重叠的server-id范围MySQL会判定为重复从库把旧连接踢掉另一种是源库大事务产生超大binlog事件超过max_allowed_packet限制。解决给每个任务分配独立且不重叠的server-id段段的宽度跟并行度匹配大事务场景把max_allowed_packet调大同时打开binlog_row_imageMINIMAL只记录变更后的镜像。MySQL事务处理繁忙的业务表尽量单独拆任务不要让一个任务同时扛几十张大表。6. 进阶一点同步链路交付前加一道数据对账与延迟巡检6.1 简单对账count与max(update_time)的双端比对中间件跑通只是过程指标真正让人敢上线的是能证明“两边数据在给定窗口内一致”。我每次交付同步链路都会在目标库侧留一张对账任务定时执行下面这类SQLSELECT COUNT(*) AS cnt, MAX(update_time) AS max_ts FROM shop.t_order; SELECT COUNT(*) AS cnt, MAX(update_time) AS max_ts FROM dw.t_order_dwd;两边count一致、最大更新时间一致说明至少没有漏跑和明显积压。粒度更细的按业务日期分组比对聚合值比如订单数与金额合计一旦有差异缩小到具体日期再查。这个方法不能证明行级完全一致但能快速暴露同步任务“静默死亡”的最坏情况。6.2 延迟巡检把事件时间和写入时间做成业务指标对账是事后校验延迟巡检才是实时健康度。我的做法是在同步写入时额外保留两个时间消息里的业务事件时间和目标表的写入时间然后定时取延迟最久的记录SELECT MAX(EXTRACT(EPOCH FROM (sync_time - event_time))) AS max_delay_sec FROM dw.t_order_dwd;这个指标超过阈值就报警。我早期上线同步任务只看进程在不在结果源端主从切换后任务照跑、数据停了几个小时都没发现后来养成了三个习惯补数脚本随时能跑、每天自动对账、延迟报警设成必接。数据同步这件事“能跑”和“可靠”之间隔着的就是这些验证动作。希望帮到你。本文还有配套的精品资源点击获取
返回列表