一、业务场景:为什么要拆历史和实时数据同步

很多做数据开发的朋友都遇到过这样的麻烦:公司的业务系统已经跑了好几年,攒了几百万甚至上亿条旧数据,现在要把这些数据同步到数据仓库做分析,同时还要保证新产生的业务数据能立刻同步过去——既要补全“旧账”,又要跟上“新账”,要是只靠一种工具,往往会出问题。 比如你用专门做实时同步的工具,跑历史数据的时候,要么速度慢得像蜗牛,要么会把业务系统的数据库压垮;要是用专门做批量同步的工具,又没法实时捕捉新数据的变化,导致分析结果滞后好几个小时甚至一天。 这时候就需要把历史同步和实时同步分开做:用适合批量处理的工具跑历史回溯,用适合实时处理的工具做增量同步,而DataX和Flink CDC就是一对非常好用的组合。

二、核心技术:DataX和Flink CDC分别擅长什么

要把这两个工具用明白,得先搞清楚它们各自的定位,就像盖房子,一个是搬旧砖的卡车,一个是运新砖的传送带,分工不同但配合默契。

2.1 DataX:批量同步的“搬运工”

DataX是阿里开源的一个离线数据同步工具,它的核心就是做“批量数据的搬运”——能支持几十种数据源之间的同步,比如从MySQL到Hive、从Oracle到MySQL、从CSV到Elasticsearch等等。 它的工作原理很简单,就是把同步任务拆成“读”和“写”两部分:读插件负责从源数据库把数据批量读出来,写插件负责把数据批量写到目标数据库,中间通过内存缓冲区做中转,速度快、稳定性高,而且能通过配置控制同步的速度、并发数,不会给源数据库造成太大压力。 举个例子,要是你有1000万条历史数据,用DataX可以分批次同步,每批同步10万条,每次同步间隔1分钟,既能把所有数据都搬过去,又不会把业务系统的MySQL压得没法用。

2.2 Flink CDC:实时同步的“传送带”

Flink CDC是基于Flink的一个增量数据捕获工具,它的核心是“实时捕捉数据的变化”——通过读取数据库的binlog(二进制日志),就能实时获取数据的插入、更新、删除操作,然后把这些变化实时同步到目标数据库。 它的工作原理和DataX完全不同:DataX是“批量拉取”,Flink CDC是“实时监听”,不需要重复扫描全表,只要监听binlog的变化就行,所以延迟非常低,一般能做到秒级甚至毫秒级的同步。 比如业务系统新增了一条订单数据,Flink CDC能在几秒钟内就把这条数据同步到数据仓库,保证分析结果的实时性。

三、协同实战:用DataX做历史回溯,用Flink CDC做实时增量

接下来我们通过一个完整的示例,一步步教大家怎么把这两个工具配合起来用,这里我们统一使用的技术栈是:源数据库MySQL 8.0、目标数据库Hive 3.1.2、DataX 3.0、Flink 1.17 + Flink CDC 2.4。

3.1 准备工作:配置源数据库的权限

不管用哪个工具,首先得给它们开权限,不然连数据库都连不上,更别说同步数据了。 首先登录MySQL,创建一个专门用于同步的用户,给这个用户分配读权限,同时要开启binlog(Flink CDC需要读取binlog):

-- 创建同步用户,密码设置为Sync@123456
CREATE USER 'sync_user'@'%' IDENTIFIED BY 'Sync@123456';
-- 给用户分配源数据库的读权限,这里源数据库是business_db,表是order表
GRANT SELECT ON business_db.order TO 'sync_user'@'%';
-- 给用户分配binlog的读权限(Flink CDC需要)
GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'sync_user'@'%';
-- 刷新权限
FLUSH PRIVILEGES;

然后修改MySQL的配置文件my.cnf,开启binlog并设置相关参数:

# 开启binlog
log-bin=mysql-bin
# 设置binlog格式为ROW(必须设置,Flink CDC才能正常读取)
binlog-format=ROW
# 设置server_id,每个MySQL实例的server_id必须唯一
server_id=1

修改完配置后,重启MySQL服务,然后登录MySQL,执行show variables like 'log_bin';,如果结果是ON,说明binlog开启成功。

3.2 第一步:用DataX做历史回溯

历史回溯的核心是“把业务系统中已经存在的旧数据,批量同步到目标数据库”,这里我们要把MySQL中business_db库的order表(假设表结构是:order_id INT、user_id INT、amount DECIMAL(10,2)、create_time DATETIME)同步到Hive的ods_order表(Hive表结构和MySQL一致)。 首先编写DataX的同步配置文件,这个文件是JSON格式,用来告诉DataX从哪里读、写到哪里、怎么读怎么写:

{
  "job": {
    "content": [
      {
        "reader": {
          "name": "mysqlreader", // 读插件:从MySQL读数据
          "parameter": {
            "username": "sync_user", // 刚才创建的同步用户名
            "password": "Sync@123456", // 同步用户密码
            "column": ["order_id", "user_id", "amount", "create_time"], // 要同步的字段
            "splitPk": "order_id", // 分片键:按order_id拆分同步任务,提高速度
            "connection": [
              {
                "table": ["order"], // 要同步的表
                "jdbcUrl": ["jdbc:mysql://192.168.1.100:3306/business_db?useUnicode=true&characterEncoding=utf8&serverTimezone=GMT%2B8"] // 源数据库连接地址
              }
            ]
          }
        },
        "writer": {
          "name": "hivewriter", // 写插件:写到Hive
          "parameter": {
            "username": "hive", // Hive的用户名
            "password": "hive", // Hive的密码
            "column": ["order_id", "user_id", "amount", "create_time"], // 要写的字段,和读的字段对应
            "preSql": ["truncate table ods_order"], // 同步前先清空Hive的表,避免重复数据
            "connection": [
              {
                "jdbcUrl": "jdbc:hive2://192.168.1.101:10000/default", // Hive的连接地址
                "table": ["ods_order"] // 要写到的Hive表
              }
            ]
          }
        }
      }
    ],
    "setting": {
      "speed": {
        "channel": 5, // 同步并发数,这里设置为5,根据服务器性能调整
        "byte": "10485760" // 同步速度限制,10MB/s,避免压垮源数据库
      }
    }
  }
}

然后把这个配置文件保存为mysql2hive_history.json,放到DataX的安装目录下,执行同步命令:

# 进入DataX安装目录
cd /opt/datax/bin
# 执行同步任务
python datax.py /opt/datax/bin/mysql2hive_history.json

执行完成后,登录Hive,执行select count(*) from ods_order;,如果结果和MySQL中order表的历史数据条数一致,说明历史回溯完成。

3.3 第二步:用Flink CDC做实时增量

历史回溯完成后,接下来要做的是“实时同步业务系统新增、更新、删除的数据”,这里我们用Flink CDC把MySQL中order表的增量数据同步到Hive的ods_order表,而且要保证增量数据不会覆盖历史数据。 首先编写Flink CDC的同步任务代码,这里我们用Java语言编写,因为Flink对Java的支持最完善:

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.connectors.mysql.table.StartupOptions;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.TableDescriptor;

public class MysqlToHiveCDC {
    public static void main(String[] args) throws Exception {
        // 1. 创建流执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1); // 设置并行度,根据服务器性能调整
        StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);

        // 2. 创建MySQL CDC源表,读取order表的增量数据
        tEnv.createTable("mysql_order", TableDescriptor.forConnector("mysql-cdc")
                .schema(Schema.newBuilder()
                        .column("order_id", DataTypes.INT())
                        .column("user_id", DataTypes.INT())
                        .column("amount", DataTypes.DECIMAL(10, 2))
                        .column("create_time", DataTypes.TIMESTAMP(3))
                        .build())
                .option("hostname", "192.168.1.100") // MySQL地址
                .option("port", "3306") // MySQL端口
                .option("username", "sync_user") // 同步用户名
                .option("password", "Sync@123456") // 同步用户密码
                .option("database-name", "business_db") // 源数据库名
                .option("table-name", "order") // 源表名
                .option("server-id", "2") // 给Flink CDC分配的server_id,要和MySQL的server_id不同
                .option("scan.startup.mode", "latest-offset") // 启动模式:从最新的binlog开始读,避免重复同步历史数据
                .build());

        // 3. 创建Hive目标表,接收增量数据
        tEnv.createTable("hive_order", TableDescriptor.forConnector("hive")
                .schema(Schema.newBuilder()
                        .column("order_id", DataTypes.INT())
                        .column("user_id", DataTypes.INT())
                        .column("amount", DataTypes.DECIMAL(10, 2))
                        .column("create_time", DataTypes.TIMESTAMP(3))
                        .build())
                .option("hive-database", "default") // Hive库名
                .option("hive-table", "ods_order") // Hive表名
                .option("sink.partition-commit.trigger", "partition-time") // 分区提交触发方式
                .option("sink.partition-commit.delay", "1h") // 分区提交延迟,避免小文件过多
                .option("sink.partition-commit.policy.kind", "metastore,success-file") // 分区提交策略
                .build());

        // 4. 把MySQL源表的数据插入到Hive目标表
        tEnv.executeSql("INSERT INTO hive_order SELECT * FROM mysql_order");

        // 5. 执行任务
        env.execute("MysqlToHiveCDC");
    }
}

然后把这段代码打包成jar包,放到Flink的安装目录下,执行同步命令:

# 进入Flink安装目录
cd /opt/flink/bin
# 提交Flink CDC任务
./flink run -c com.example.MysqlToHiveCDC /opt/flink/jar/mysql2hive-cdc.jar

提交完成后,登录Flink的WebUI(一般是http://192.168.1.102:8081),就能看到任务在运行,这时候业务系统新增的订单数据,就能实时同步到Hive的ods_order表了。

3.4 协同的关键:避免重复数据

很多朋友在配合的时候会遇到重复数据的问题,比如历史回溯的时候,Flink CDC也在同步,导致同一条数据出现两次。要解决这个问题,核心是控制两个任务的启动时间和启动模式:

  1. 先启动DataX的历史回溯任务,等历史回溯完全完成后,再启动Flink CDC的实时增量任务;
  2. Flink CDC的启动模式要设置为latest-offset,也就是从最新的binlog开始读,这样就不会重复读取已经被DataX同步过的历史数据;
  3. 如果历史回溯需要很长时间,比如几个小时甚至几天,这期间业务系统会产生新数据,Flink CDC可以先启动,但要设置为initial模式(从表的开头读),然后等DataX完成后,再把Flink CDC的启动模式改成latest-offset,不过这种方式比较复杂,一般建议先跑历史回溯,再跑实时增量。

四、技术优缺点分析

4.1 DataX的优缺点

优点:

  1. 支持的数据源多,能覆盖大部分常用的数据库、文件系统、大数据组件;
  2. 批量同步速度快,能通过并发、分片等方式优化同步效率,适合处理大量历史数据;
  3. 配置简单,不需要编写代码,只要写JSON配置文件就能完成同步;
  4. 稳定性高,阿里开源后经过了大量的线上验证,很少出现数据丢失的情况。 缺点:
  5. 只能做离线批量同步,不能实时捕捉数据变化,无法满足实时同步的需求;
  6. 同步延迟高,每次同步都要全表扫描,适合T+1的同步场景,不适合T+0的场景;
  7. 没有断点续传功能,要是同步过程中中断,只能重新开始同步。

4.2 Flink CDC的优缺点

优点:

  1. 实时性高,延迟能做到秒级甚至毫秒级,适合实时增量同步的场景;
  2. 支持断点续传,要是同步过程中中断,下次启动会从上次中断的地方继续,不会重复同步;
  3. 支持数据的插入、更新、删除操作,能完整捕捉数据的变化;
  4. 能和Flink的其他功能结合,比如流批处理、复杂事件处理等,适合复杂的实时数据处理场景。 缺点:
  5. 不适合处理大量历史数据,要是用Flink CDC同步历史数据,速度会非常慢,而且会给源数据库造成很大压力;
  6. 配置复杂,需要编写代码或者复杂的SQL,对新手不太友好;
  7. 依赖binlog,要是源数据库的binlog配置不正确,Flink CDC就无法正常工作。

4.3 协同的优缺点

优点:

  1. 兼顾了历史数据和实时数据的同步,既能补全旧数据,又能跟上新数据的变化;
  2. 充分发挥了两个工具的优势,避免了各自的缺点,提高了同步效率和稳定性;
  3. 能满足大部分业务场景的需求,比如数据仓库的同步、数据备份、数据迁移等。 缺点:
  4. 需要同时维护两个工具,增加了运维的复杂度;
  5. 要是两个工具的配置不正确,很容易出现重复数据或者数据丢失的问题;
  6. 没有统一的管理界面,需要分别监控两个工具的运行状态。

五、注意事项

  1. 源数据库的权限配置:不管是DataX还是Flink CDC,都需要给同步用户分配足够的权限,同时要避免分配过高的权限,比如root权限,避免安全风险;
  2. 源数据库的binlog配置:Flink CDC需要读取binlog,所以必须把binlog的格式设置为ROW,server_id必须唯一,同时要保证binlog的保留时间足够长,要是binlog被删除了,Flink CDC就无法正常同步;
  3. 同步速度的控制:DataX同步历史数据的时候,要根据源数据库的性能设置并发数和同步速度,避免压垮源数据库;Flink CDC同步的时候,要根据服务器的性能设置并行度,避免资源浪费;
  4. 重复数据的处理:要严格控制两个任务的启动时间和启动模式,避免重复同步数据;要是出现重复数据,可以通过主键去重的方式解决;
  5. 监控和告警:要分别监控DataX和Flink CDC的运行状态,要是出现任务中断、数据延迟过高的情况,要及时告警,避免影响业务;
  6. 数据一致性的保证:要定期检查源数据库和目标数据库的数据是否一致,比如通过count(*)的方式对比,要是出现不一致的情况,要及时排查原因。

六、文章总结

DataX和Flink CDC的协同,是一种非常实用的数据同步方案,它通过拆分历史同步和实时同步的任务,充分发挥了两个工具的优势,既能高效、稳定地完成历史数据的回溯,又能实时捕捉增量数据的变化,满足了大部分业务场景的需求。 在实际使用的时候,要根据自己的业务场景选择合适的配置,严格遵守注意事项,避免出现重复数据、数据丢失、源数据库压力过大等问题。同时,要做好监控和运维工作,保证同步任务的稳定运行。