一年多前写的用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
   );
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐