spark 实现sqoop类似功能的数据同步工具
·
一年多前写的用spark 实现sqoop类似功能的数据同步工具,拿出来分享一下
优点:
1.用配置代替sqoop命令,维护简单,添加新增同步表或下线同步表,只需在sql中增加或删除即可
2.实测性能比sqoop好,当然前提是内存资源足够
实现功能:
1.通过配置需要同步的表清单元数据,实现多表并行从mysql同步到hive
2.通过配置表同步类型,实现全量同步,增量同步,是否重跑功能
2.通过配置表的并行分区字段,分区度(并行度),实现表数据分区并行同步
3.同步过程中,动态检测源表表结构字段变化,并做相应处理,在源表字段有增减的情况,同步能顺利完成
主逻辑代码:
package com.hn.data
import java.text.SimpleDateFormat
import java.util
import java.util.{Calendar, Date, Locale}
import org.apache.log4j.{Level, Logger}
import scala.collection.JavaConversions._
import scala.collection.parallel.ForkJoinTaskSupport
import scala.util.{Failure, Success, Try}
object newDataSyncMysql2Hive {
Logger.getLogger("org").setLevel(Level.ERROR)
val logger = Logger.getLogger("newDataSyncMysql2Hive")
logger.setLevel(Level.ERROR)
//构建查询源msyql表的where过滤条件语句
def getMysqlFilterStatment(syncType: String, mysqlFilterCol: Array[String], processDate: String, processDateOneMoreDay: String): String = {
var sql = ""
//增量同步
if (syncType.equals("inc")) {
for (i <- mysqlFilterCol.indices) {
if (i == mysqlFilterCol.length - 1) {
val col = mysqlFilterCol(i)
sql = sql + s"($col >='$processDate' and $col < '$processDateOneMoreDay') "
} else {
val col = mysqlFilterCol(i)
sql = sql + s"($col >='$processDate' and $col < '$processDateOneMoreDay') or "
}
}
}
//全量同步
else {
sql = sql + " 1 = 1"
}
sql
}
def main(args: Array[String]): Unit = {
//数据同步的日期,作为hive表的日期分区
val processDate = args(0) //"20181219"
//需要同步的表清单的并行度,如清单中有10张表,若设置并行为2,则把表清单分成2分,每份5张表,两份并行通读
val parallelCnt = args(1) //"1"
//重跑时设置表是否强制重跑 1:不管synced表,强制重新同步,0:从synced表判断是否重新同步
val ifForceReSync = args(2) //"1"
//根据业务,可能有些表需要晚上同步,有些表需要白天同步,把这些设置批次
val importBatchNum = args(3).toInt //每天第几批运行
//mysql url
val metaMysqlUrl = ""
//用户名
val metaUserName = ""
// 密码
val metaPassword = ""
// data_sync_meta.meta_data_sync_mysql2hive 存储表清单元数据,存放在mysql表
val syncTbTableSql = s"select * from data_sync_meta.meta_data_sync_mysql2hive where importBatchNum = $importBatchNum"
// data_sync_meta.meta_data_synced_mysql2hive用来存储当前同步完成的表,存放在hive
val syncedTbTable = "data_sync_meta.meta_data_synced_mysql2hive"
// 格式化日期
val dataProcessStartDate = Utils.dateFormat(processDate)
// 获取spark session
val spark = Utils.getSparkSession("newDataSyncMysql2Hive1")
// 获取同步表清单
val tableListDF = Utils.getMysqlJdbcDFSingle(spark, metaMysqlUrl, syncTbTableSql, metaUserName, metaPassword)
// 将tableListDF collect成scala原生并行数组。ps:使用df.foreach并行报NPE错误
val tableList = tableListDF.collect().par
// 调节表间并行度
tableList.tasksupport = new ForkJoinTaskSupport(new scala.concurrent.forkjoin.ForkJoinPool(parallelCnt.toInt))
//spark.sparkContext.textFile(tableListFile).repartition(parallelCnt.trim.toInt).filter(row => row.startsWith("mysqlTable") != true)
import spark.implicits._
//获取已经同步完成的表清单
val syncedTableList = spark.sql(s"select synced_table from $syncedTbTable where day = '$processDate'").map(row => row.getString(0)).collect()
//tableListDF.collect().foreach(println)
val df = new SimpleDateFormat("yyyy-MM-dd", Locale.getDefault)
val df1 = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
val calendar = Calendar.getInstance(Locale.getDefault)
calendar.setTime(df.parse(dataProcessStartDate))
//处理日期加1天,作为处理end time
calendar.add(Calendar.DATE, 1)
val dataProcessEndDate = df.format(calendar.getTime)
val processStartTime = System.currentTimeMillis()
//这里要将spark broadcast一下,否则rdd里面要用spark报NPE错误
val broadCastSpark = spark.sparkContext.broadcast(spark)
// 处理表清单,获取每张表的相关同步配置
tableList.foreach { row =>
Try {
val syncTableSeq = row.getLong(0) //表序号
val mysqlTable = row.getString(1) //mysql源表名
val hiveTable = row.getString(2) //hive目标表名
val mysqlDataFilterCols = row.getString(3) //源表过滤时间字段
val sparkTmpTbName = row.getString(4) //spark临时表名
val mysqlUrl = row.getString(5) //mysqlurl
val userName = row.getString(6) //mysql用户名
val pwd = row.getString(7) // mysql pwd
val isHivePartitionTable = row.getInt(8) //目标hive表是否是分区表
val hivePartitionCol = row.getString(9) // hive表分区字段
val syncType = row.getString(10) // 同步类型,全量 or 增量
val ifTableParallel = row.getInt(11) // 单表是否并行导入
val runTimePartitionColumn = row.getString(12) // 多并行同步表时,用来计算并行度(分区)的字段
val numPartitions = row.getInt(13) // 表并行度
if (syncedTableList.contains(hiveTable) && ifForceReSync.equals("0")) {
println(s"===========>$hiveTable 已经同步,此次忽略")
} else {
println(s"===========>第 $syncTableSeq 个表 $hiveTable 开始同步")
// 查询mysql数据源
val mysqlFilterCol = mysqlDataFilterCols.split("&")
// 构建查询mysql源表的where过滤语句
val mysqlFilterStatment = getMysqlFilterStatment(syncType, mysqlFilterCol, dataProcessStartDate, dataProcessEndDate)
// 构建查询mysql源表的完整sql语句
val mysqlResSql = s"select * from $mysqlTable where $mysqlFilterStatment"
println(s">>>>> 查询源表sql: $mysqlResSql")
//println(broadCastSpark.value, mysqlUrl, sql, userName, pwd)
// 获取mysql表最大最小id,spark在根据这个两个值,和配置的表同步并行,并行拉取mysql数据
// 如最小id是1,最大id是100,并行度设为10,则spark按并行度分为10个分区,每个分区10条数据
val mysqlDataDF = ifTableParallel match {
// 如果表设置并行导入
case 1 => {
val tableMinMaxIdSql = s"select min($runTimePartitionColumn) as minid, max($runTimePartitionColumn) as maxid from $mysqlTable where $mysqlFilterStatment "
println(s">>>>> 查询源表最大最小id sql: $tableMinMaxIdSql")
Utils.getMysqlJdbcDFSingle(broadCastSpark.value, mysqlUrl, tableMinMaxIdSql, userName, pwd).registerTempTable("tableMinMaxId")
val tableMinMaxId = spark.sql("select cast(coalesce(minid,0) as bigint) as minid,cast(coalesce(maxid,999999999) as bigint) as maxid from tableMinMaxId")
// tableMinMaxId.printSchema()
// tableMinMaxId.show(10,false)
val minId = tableMinMaxId.map(row => row.getLong(0)).take(1)(0)
val maxId = tableMinMaxId.map(row => row.getLong(1)).take(1)(0)
// spark在根据这个两个值,和配置的表同步并行,并行拉取mysql数据
Utils.getMysqlJdbcDFParallel(broadCastSpark.value, mysqlUrl, mysqlResSql, userName, pwd, runTimePartitionColumn, minId, maxId, numPartitions)
}
// 如果表设置非并行导入,则按一个并行
case 0 => Utils.getMysqlJdbcDFSingle(broadCastSpark.value, mysqlUrl, mysqlResSql, userName, pwd)
}
//spark 临时表
mysqlDataDF.registerTempTable(sparkTmpTbName)
// 获取spark临时表schema
val mysqlSchema = broadCastSpark.value.sql(s"desc $sparkTmpTbName").select($"col_name").map(row => row.getString(0).toLowerCase).collect
import spark.implicits._
//获取hive目标表字段列表
val hiveSchema = broadCastSpark.value.sql(s"desc $hiveTable").select($"col_name").filter(!$"col_name".startsWith("#") && $"col_name" =!= hivePartitionCol)
.map(row => row.getString(0)).collect()
//源表字段列表与目标表字段列表对比,如果源表字段多余目标,忽略,可以线下增加字段与源表保持一致
//如果源表字段少于目标表,添加null as 补充字段,以免报找不到错误
val tgtCols = new util.ArrayList[String]()
for (h <- hiveSchema) {
if (mysqlSchema.contains(h.toString.toLowerCase)) {
// 去除字段中回车换行
tgtCols.add(s"regexp_replace(" + h.toString.toLowerCase + s", '\\n|\\t|\\r', '') as " + h.toString.toLowerCase)
}
else {
tgtCols.add(s"null as $h")
}
}
//构建插入目标表的字段列表
val hiveCol = tgtCols.mkString(",")
if (isHivePartitionTable == 1) {
//删除目标表目标分区
val dropPartitionSql = s"alter table $hiveTable drop if exists partition($hivePartitionCol='$processDate')"
println(s">>>>> drop语句 :$dropPartitionSql")
broadCastSpark.value.sql(dropPartitionSql)
//同步数据到目标表,目标分区
val insertSql = s"insert into table $hiveTable partition($hivePartitionCol='$processDate')" +
s" select $hiveCol from $sparkTmpTbName"
println(s">>>>> insert语句 :$insertSql")
broadCastSpark.value.sql(insertSql)
//将同步完成的表名插入已同步表清单表中
val EtlCreateTime = df1.format(new Date)
val syncedInsertSql = s"insert into table $syncedTbTable partition(day = '$processDate') values ('$hiveTable','$EtlCreateTime')"
println(s">>>>> synced insert语句 :$syncedInsertSql")
broadCastSpark.value.sql(syncedInsertSql)
println(s"===========>$hiveTable 同步完成")
} else {
val insertSql = s"insert overwrite table $hiveTable " +
s" select $hiveCol from $sparkTmpTbName"
println(s">>>>> insert语句 :$insertSql")
broadCastSpark.value.sql(insertSql)
val EtlCreateTime = df1.format(new Date)
broadCastSpark.value.sql(s"insert into table $syncedTbTable partition(day = '$processDate') values ('$hiveTable','$EtlCreateTime')")
println(s"===========>$hiveTable 同步完成")
}
}
} match {
case Success(s) =>
case Failure(e) => DingTalkWarn.sendWarn("数据同步:" + row.getString(1) + "-" + row.getString(2))
println(e)
}
}
spark.stop()
val processEndTime = System.currentTimeMillis()
println(s">>>>>>>> 总计运行 " + (processEndTime - processStartTime) / 1000 + " s")
}
}
Utils.cala
package com.hn.data
import java.text.SimpleDateFormat
import java.util._
import org.apache.spark.sql.{SparkSession, _}
object Utils {
val sparkSqlWarehouseDir = ""
val hbaseZookeeperQuorum = ""
def getSparkSession(appName: String): SparkSession = {
val submitType = System.getProperty("os.name").startsWith("Windows") match {
case true => "local[4]"
case _ => "yarn"
}
val spark = submitType match {
case "local[4]" => SparkSession.builder()
.master(submitType)
.appName(appName)
.config("spark.sql.warehouse.dir", sparkSqlWarehouseDir)
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config("spark.executor.memory", "12g")
.config("spark.driver.memory", "2g")
.config("spark.driver.cores", "4")
.config("spark.debug.maxToStringFields", "100")
.config("spark.kryoserializer.buffer.max", "128")
.config("spark.debug.maxToStringFields", "100")
.config("hbase.zookeeper.quorum", hbaseZookeeperQuorum)
// .config("spark.hbase.host", sparkHbaseHost)
.enableHiveSupport()
.getOrCreate()
case _ => SparkSession.builder()
.master(submitType)
.appName(appName)
.config("spark.sql.warehouse.dir", sparkSqlWarehouseDir)
.config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
.config("spark.debug.maxToStringFields", "100")
.config("spark.kryoserializer.buffer.max", "128")
.config("spark.debug.maxToStringFields", "100")
.config("hbase.zookeeper.quorum", hbaseZookeeperQuorum)
// .config("spark.hbase.host", sparkHbaseHost)
.enableHiveSupport()
.getOrCreate()
}
spark
}
def getMysqlJdbcDFParallel(spark: SparkSession, url: String, querySql: String, userName: String, passWord: String, runTimePartitionColumn: String, lowerBound: Long, upperBound: Long, numPartitions: Int): DataFrame = {
val readConnProperties = new Properties()
readConnProperties.put("driver", "com.mysql.jdbc.Driver")
readConnProperties.put("user", userName)
readConnProperties.put("password", passWord)
//readConnProperties.put("fetchsize", "1000")
spark.read.jdbc(
url,
s"($querySql) t", // 注意括号和表别名,必须得有,这里可以过滤数据
runTimePartitionColumn, //mysql并发分区字段
lowerBound, //并发字段最小值
upperBound, //并发字段最大值
numPartitions, //并发值
readConnProperties)
}
def getMysqlJdbcDFSingle(spark: SparkSession, url: String, querySql: String, userName: String, passWord: String): DataFrame = {
val readConnProperties = new Properties()
readConnProperties.put("driver", "com.mysql.jdbc.Driver")
readConnProperties.put("user", userName)
readConnProperties.put("password", passWord)
//readConnProperties.put("fetchsize", "1000")
spark.read.jdbc(
url,
s"($querySql) t", // 注意括号和表别名,必须得有,这里可以过滤数据
readConnProperties)
}
def dateFormat(dateStr: String): String = {
val df = new SimpleDateFormat("yyyyMMdd", Locale.getDefault)
val df1 = new SimpleDateFormat("yyyy-MM-dd", Locale.getDefault)
dateStr.contains("-") match {
case false =>
df.setLenient(false)
df1.format(df.parse(dateStr))
case true =>
df1.setLenient(false)
df.format(df1.parse(dateStr))
}
}
}
表清单元数据存储表和已同步完成清单表
-- 同步表清单元数据表(mysql)
CREATE TABLE `data_sync_meta`.`meta_data_sync_mysql2hive` (
`id` bigint(20) NOT NULL AUTO_INCREMENT,
`mysqlTable` varchar(100) DEFAULT NULL COMMENT 'mysql表',
`hiveTable` varchar(100) DEFAULT NULL COMMENT 'hive表',
`mysqlDataFilterCols` varchar(100) DEFAULT NULL COMMENT '增量同步时mysql时间戳过滤字段,有多个时用&隔开,full表填nd',
`sparkTmpTbName` varchar(100) DEFAULT NULL COMMENT 'spark临时表',
`mysqlUrl` varchar(255) DEFAULT NULL COMMENT 'mysql jdbc url',
`userName` varchar(100) DEFAULT NULL COMMENT 'username',
`pwd` varchar(100) DEFAULT NULL COMMENT 'pwd',
`isHivePartitionTable` int(11) DEFAULT NULL COMMENT 'hive表是否分区表 1是 0否',
`hivePartitionCol` varchar(100) DEFAULT NULL COMMENT ' hive表分区字段名',
`syncType` varchar(20) DEFAULT NULL COMMENT '同步类型 全量 full 增量inc',
`ifTableParallel` int(11) DEFAULT NULL COMMENT '表是否切割并行导入 1是,0否',
`runTimePartitionColumn` varchar(100) DEFAULT NULL COMMENT ' 并行导入时,切割字段名,必须int或相关数字类型',
`numPartitions` int(11) DEFAULT NULL COMMENT '并行度',
`importBatchNum` int(11) DEFAULT NULL COMMENT '当天第几批次导入',
`create_time` datetime DEFAULT NULL,
`update_time` datetime DEFAULT NULL,
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
--同步完成清单表(hive)
create table data_sync_meta.meta_data_synced_mysql2hive(
syncedtable String comment "已经同步完成的表",
etl_time String comment "etl时间"
) partitioned by (day string)
配置表清单示例:
INSERT INTO `data_sync_meta`.`meta_data_sync_mysql2hive` (`id`, `mysqlTable`, `hiveTable`, `mysqlDataFilterCols`, `sparkTmpTbName`, `mysqlUrl`, `userName`, `pwd`, `isHivePartitionTable`, `hivePartitionCol`, `syncType`, `ifTableParallel`, `runTimePartitionColumn`, `numPartitions`, `importBatchNum`, `create_time`, `update_time`)
VALUES (null --id
,'xx.xxx'-- mysql表
,'ods_xxx.xxx' -- hive表
,'nd' -- 增量同步时mysql时间戳过滤字段,有多个时用&隔开,full表和source表填nd
,'xxxx' -- spark临时表
,'jdbc:mysql://x.x.x.x:3306/xxx?zeroDateTimeBehavior=convertToNull&tinyInt1isBit=false' -- mysql jdbc url
,'username' -- mysql username
,'pwd' -- mysql pwd
, 1 -- hive表是否分区表 1是 0否
, 'pday' -- hive表分区字段名
, 'full' -- 同步类型 全量 full 增量inc
, 1 -- 表是否切割并行导入 1是,0否
, 'id' -- 并行导入时,切割字段名,必须int或相关数字类型
, 5 -- 并行度
, 99 -- 当天第几批次导入
, null -- etl_create_time
, null -- etl_update_time
);
更多推荐
所有评论(0)