一、为啥要做Hive和MySQL双向同步?

很多做数据相关工作的朋友,都会遇到这么个场景:公司有两套数据系统,一套是存海量历史数据、做复杂分析的Hive,另一套是给业务用、查改快的MySQL。比如电商公司,用户的历史消费数据都存在Hive里做用户画像分析,而用户的实时余额、最新订单又存在MySQL里给前端调用。但有时候两边数据得互通:比如Hive里分析出用户是高价值用户,得把这个标签同步到MySQL,让业务系统能给高价值用户推专属活动;反过来,MySQL里的实时订单数据,又得同步到Hive,更新用户画像。

要是手动同步,不仅麻烦,还容易出错,所以就得找个能自动同步的工具,DataX就是很多人会选的工具之一。

二、先搞懂DataX是什么?

可能有人没接触过DataX,这里先简单说下,别觉得它复杂,其实就是个数据搬运工。它是阿里开源的一个数据同步工具,专门用来把不同数据库、不同存储系统之间的数据搬来搬去。它的核心逻辑很简单:先把源端的数据读出来,转成中间格式,再写到目标端。

比如要把Hive的数据搬到MySQL,DataX就先连Hive把数据读出来,转成通用的中间格式,再连MySQL把数据写进去;反过来搬也一样。而且它支持的数据源特别多,除了Hive、MySQL,还支持Oracle、MongoDB、Elasticsearch这些,基本能覆盖常用的场景。

三、准备工作:先搭好环境

要做同步,得先把需要的环境都弄好,这里说的环境都是日常用的,没有复杂的配置。

3.1 先装DataX

DataX是用Java写的,所以得先装JDK,版本至少是1.8以上,这个大家基本都有,要是没有的话,装个OpenJDK8就行。装完JDK,得配置环境变量,把JAVA_HOME设对,这个网上搜一下就有,很简单。

然后去DataX的GitHub仓库下载最新的包,地址是https://github.com/alibaba/DataX,下载完解压出来就行,不用编译,直接就能用。解压完的文件夹里,bin目录是用来启动同步任务的,plugin目录是各种数据源的插件,比如读Hive的插件、写MySQL的插件都在这。

3.2 准备Hive和MySQL的测试数据

为了演示方便,我们先建两个测试表,一个Hive的,一个MySQL的,结构要对应上,不然同步的时候会出问题。

先建MySQL的测试表,用root用户登录MySQL,执行下面的SQL:

-- 创建测试库test_db
CREATE DATABASE IF NOT EXISTS test_db;
USE test_db;
-- 创建用户标签表,用来存用户的高价值标签
CREATE TABLE IF NOT EXISTS user_tag (
    user_id INT PRIMARY KEY COMMENT '用户ID',
    tag VARCHAR(50) COMMENT '用户标签,比如高价值',
    update_time DATETIME COMMENT '更新时间'
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

再建Hive的测试表,Hive是基于HDFS的,所以得先保证Hadoop集群是正常的,然后用hive命令行建表:

-- 创建测试库test_db
CREATE DATABASE IF NOT EXISTS test_db;
USE test_db;
-- 创建用户画像表,用来分析用户的价值标签
CREATE TABLE IF NOT EXISTS user_profile (
    user_id INT COMMENT '用户ID',
    tag VARCHAR(50) COMMENT '用户标签',
    update_time TIMESTAMP COMMENT '更新时间'
) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t';

这里要注意,Hive表的字段类型和MySQL的要尽量对应,比如user_id都是INT,tag都是字符串,update_time都是时间类型,这样同步的时候才不会出类型转换的问题。

四、第一次同步:Hive到MySQL

先做从Hive到MySQL的同步,比如Hive里分析出用户的标签,要同步到MySQL给业务用。

4.1 写DataX的同步配置

DataX的同步任务是靠JSON配置文件来定义的,这个配置文件就相当于给搬运工的指令,告诉它从哪搬、搬什么、搬到哪。

这里用的技术栈是DataX + Hive + MySQL,配置文件的具体内容和注释如下:

{
    "job": {
        "content": [
            {
                "reader": {
                    "name": "hdfsreader", // 读Hive的插件,Hive的数据存在HDFS上,所以用hdfsreader
                    "parameter": {
                        "path": "/user/hive/warehouse/test_db.db/user_profile", // Hive表对应的HDFS路径,这个可以在Hive里用desc formatted user_profile看
                        "defaultFS": "hdfs://hadoop-master:9000", // HDFS的地址,换成自己集群的地址
                        "column": ["user_id", "tag", "update_time"], // 要同步的字段,和Hive表的字段对应
                        "fieldDelimiter": "\t", // Hive表的字段分隔符,和建表时的一致
                        "encoding": "UTF-8" // 编码,避免中文乱码
                    }
                },
                "writer": {
                    "name": "mysqlwriter", // 写MySQL的插件
                    "parameter": {
                        "writeMode": "update", // 写入模式,这里用update,要是主键重复就更新,不重复就插入
                        "username": "root", // MySQL的用户名,换成自己的
                        "password": "123456", // MySQL的密码,换成自己的
                        "column": ["user_id", "tag", "update_time"], // 要写入的字段,和MySQL表的字段对应
                        "preSql": ["SET NAMES utf8mb4;"], // 写入前执行的SQL,设置编码,避免中文乱码
                        "connection": [
                            {
                                "jdbcUrl": "jdbc:mysql://mysql-server:3306/test_db?useUnicode=true&characterEncoding=utf8&serverTimezone=GMT%2B8", // MySQL的连接地址,换成自己的,注意时区要对
                                "table": ["user_tag"] // 要写入的MySQL表名
                            }
                        ]
                    }
                }
            }
        ],
        "setting": {
            "speed": {
                "channel": 3 // 同步的线程数,根据自己的服务器性能调,一般3-5就够
            }
        }
    }
}

这里要注意几个细节:一是HDFS的路径,要是不知道自己Hive表的路径,可以在Hive命令行里执行desc formatted user_profile;,然后找Location那一行,就是对应的HDFS路径;二是MySQL的时区,要是不设置对,时间类型同步过去会差8个小时;三是writeMode,这里用update是因为user_id是主键,要是Hive里的用户标签更新了,同步到MySQL的时候会覆盖旧的标签,要是用insert的话,主键重复就会报错。

4.2 启动同步任务

配置文件写好后,保存成hive2mysql.json,然后用DataX的bin目录下的datax.py来启动任务,命令如下:

# 进入DataX的bin目录,然后执行下面的命令,换成自己的配置文件路径
python /datax/bin/datax.py /datax/config/hive2mysql.json

执行完后,要是没有报错,就可以去MySQL里查user_tag表,看数据是不是同步过来了。比如在Hive的user_profile表里插一条测试数据:

-- Hive里插测试数据
INSERT INTO TABLE user_profile VALUES (1001, '高价值用户', '2024-05-20 12:00:00');

然后再执行同步任务,去MySQL里查:

-- MySQL里查
SELECT * FROM user_tag;

就能看到user_id为1001的标签同步过来了。

五、第二次同步:MySQL到Hive

接下来做从MySQL到Hive的同步,比如MySQL里的实时订单数据,要同步到Hive更新用户画像。

5.1 写MySQL到Hive的同步配置

同样是写JSON配置文件,技术栈还是DataX + Hive + MySQL,具体内容和注释如下:

{
    "job": {
        "content": [
            {
                "reader": {
                    "name": "mysqlreader", // 读MySQL的插件
                    "parameter": {
                        "username": "root", // MySQL的用户名
                        "password": "123456", // MySQL的密码
                        "column": ["user_id", "tag", "update_time"], // 要同步的字段
                        "splitPk": "user_id", // 用来分块同步的字段,必须是数字类型,比如INT,这样能提高同步速度
                        "connection": [
                            {
                                "jdbcUrl": "jdbc:mysql://mysql-server:3306/test_db?useUnicode=true&characterEncoding=utf8&serverTimezone=GMT%2B8", // MySQL的连接地址
                                "table": ["user_tag"] // 要读的MySQL表名
                            }
                        ]
                    }
                },
                "writer": {
                    "name": "hdfswriter", // 写Hive的插件,Hive的数据存在HDFS上,所以用hdfswriter
                    "parameter": {
                        "defaultFS": "hdfs://hadoop-master:9000", // HDFS的地址
                        "path": "/user/hive/warehouse/test_db.db/user_profile", // Hive表对应的HDFS路径
                        "column": [
                            {"name": "user_id", "type": "INT"}, // 字段名和类型,和Hive表的对应
                            {"name": "tag", "type": "STRING"},
                            {"name": "update_time", "type": "TIMESTAMP"}
                        ],
                        "fieldDelimiter": "\t", // 字段分隔符,和Hive表的一致
                        "writeMode": "append", // 写入模式,这里用append,追加数据,要是Hive里有旧数据,会保留,新数据加在后面
                        "encoding": "UTF-8" // 编码
                    }
                }
            }
        ],
        "setting": {
            "speed": {
                "channel": 3 // 同步的线程数
            }
        }
    }
}

这里要注意splitPk,这个字段是用来分块同步的,要是MySQL里的数据很多,比如几百万条,splitPk能把数据分成好几块,每个线程同步一块,速度会快很多,但是splitPk必须是数字类型,比如INT、BIGINT,要是用字符串类型的字段,会报错。还有writeMode,这里用append,因为Hive里的用户画像数据是累计的,新同步过来的数据要追加,要是用overwrite的话,会把Hive里原来的数据都覆盖掉,一定要注意。

5.2 启动同步任务

配置文件保存成mysql2hive.json,然后启动任务:

# 启动MySQL到Hive的同步任务
python /datax/bin/datax.py /datax/config/mysql2hive.json

执行完后,去Hive里查user_profile表,看数据是不是同步过来了。比如在MySQL的user_tag表里插一条测试数据:

-- MySQL里插测试数据
INSERT INTO user_tag VALUES (1002, '潜力用户', '2024-05-20 13:00:00');

然后再执行同步任务,去Hive里查:

-- Hive里查
SELECT * FROM user_profile;

就能看到user_id为1002的标签同步过来了。

六、怎么实现自动同步?

上面的同步都是手动执行的,要是每次都手动跑,还是很麻烦,所以得做自动同步。常见的自动同步方式有两种,一种是用Linux的定时任务crontab,另一种是用调度平台比如XXL-Job、Airflow。

6.1 用crontab做定时同步

crontab是Linux自带的定时任务工具,用起来很简单。比如我们想每天凌晨2点同步一次Hive到MySQL,每天凌晨3点同步一次MySQL到Hive,就可以这么做。

首先写两个脚本,一个是同步Hive到MySQL的脚本,保存成hive2mysql.sh:

#!/bin/bash
# 每天凌晨2点同步Hive到MySQL
python /datax/bin/datax.py /datax/config/hive2mysql.json

另一个是同步MySQL到Hive的脚本,保存成mysql2hive.sh:

#!/bin/bash
# 每天凌晨3点同步MySQL到Hive
python /datax/bin/datax.py /datax/config/mysql2hive.json

然后给两个脚本加执行权限:

chmod +x /datax/scripts/hive2mysql.sh
chmod +x /datax/scripts/mysql2hive.sh

最后编辑crontab:

crontab -e

在打开的文件里加下面两行:

# 每天凌晨2点执行Hive到MySQL的同步
0 2 * * * /datax/scripts/hive2mysql.sh
# 每天凌晨3点执行MySQL到Hive的同步
0 3 * * * /datax/scripts/mysql2hive.sh

保存退出,crontab就会每天定时执行这两个任务了。要是想更灵活的调度,比如每小时同步一次,或者只在工作日同步,也可以改crontab的时间规则,这个网上搜一下crontab的语法就懂了。

6.2 同步的增量更新

上面的同步是全量同步,就是每次同步都会把所有数据都搬一遍,要是数据量很大,比如几千万条,全量同步会很慢,还会占用很多资源。所以一般实际用的时候,都是做增量同步,就是只同步上次同步之后新增或者更新的数据。

增量同步的核心是加一个时间字段,比如update_time,每次同步的时候,只同步update_time大于上次同步时间的数据。比如我们要做MySQL到Hive的增量同步,就可以在配置文件的reader里加一个where条件:

"where": "update_time > '${last_sync_time}'"

这里的${last_sync_time}是一个变量,每次同步前,先查上次同步的时间,然后把这个变量替换成具体的时间,比如上次同步的时间是2024-05-20 12:00:00,这次就只同步update_time大于这个时间的数据。

实现增量同步的具体步骤:

  1. 每次同步前,查上次同步的时间,比如存在一个日志文件里,或者存在一个MySQL的配置表里。
  2. 把这个时间替换到配置文件的where条件里。
  3. 执行同步任务。
  4. 同步完成后,把当前的时间更新为上次同步的时间,存在日志文件或者配置表里。

这样每次同步就只会搬新增或者更新的数据,速度会快很多,资源占用也少。

七、应用场景、优缺点和注意事项

7.1 应用场景

DataX实现Hive和MySQL双向同步的应用场景很多,除了开头说的用户标签同步,还有:

  • 电商的实时订单同步:MySQL里的实时订单同步到Hive,做订单分析;Hive里的订单分析结果(比如热销商品)同步到MySQL,给前端推荐用。
  • 金融的风险数据同步:Hive里分析出的风险用户标签同步到MySQL,给风控系统用;MySQL里的实时交易数据同步到Hive,更新风险模型。
  • 内容平台的内容标签同步:Hive里分析出的内容标签同步到MySQL,给内容推荐系统用;MySQL里的内容实时点击数据同步到Hive,更新内容模型。

7.2 技术优缺点

优点:

  • 开源免费,不用花钱,社区也比较活跃,遇到问题能找到解决方案。
  • 支持的数据源多,除了Hive和MySQL,还能支持其他很多数据库和存储系统,以后要是换数据源,不用换工具。
  • 配置简单,只要写JSON配置文件,不用写复杂的代码,上手快。
  • 性能不错,支持多线程同步,数据量大的时候也能跑得比较快。

缺点:

  • 全量同步的时候,要是数据量很大,速度会比较慢,增量同步需要自己实现,没有现成的功能。
  • 没有可视化界面,配置和监控都得靠命令行或者自己写脚本,对新手不太友好。
  • 同步过程中要是出问题,比如网络断了,不会自动重试,得手动重新执行任务。

7.3 注意事项

  • 字段类型要对应:Hive和MySQL的字段类型要尽量对应,比如Hive的INT对应MySQL的INT,Hive的STRING对应MySQL的VARCHAR,不然同步的时候会出类型转换错误。
  • 编码要一致:两边的编码都要设成UTF-8,避免中文乱码,MySQL的连接地址里要加useUnicode=true&characterEncoding=utf8,Hive的配置里要加encoding=UTF-8。
  • 时区要一致:MySQL的连接地址里要加serverTimezone=GMT%2B8,不然时间类型同步过去会差8个小时。
  • 主键要注意:Hive到MySQL同步的时候,要是用update模式,必须有主键,不然会重复插入数据;MySQL到Hive同步的时候,要是用append模式,要注意不要重复同步数据,最好用增量同步。
  • 资源占用:同步任务会占用服务器的CPU、内存、网络资源,所以最好在业务低峰期执行,比如凌晨,避免影响业务系统。

八、文章总结

用DataX实现Hive和MySQL的双向同步,其实就是把DataX当成一个数据搬运工,配置好从哪搬、搬什么、搬到哪,然后启动任务就行。要是想自动同步,就加个定时任务,要是数据量大,就做增量同步。

整个过程没有复杂的技术,只要把配置文件写对,注意一些细节,比如字段类型、编码、时区,就能实现稳定的双向同步。这个方法不仅适合Hive和MySQL,也适合其他很多数据源之间的同步,只要改改配置文件里的插件和参数就行。