【问题标题】:sqoop import-all-tables slow and sequence files are custom java objectssqoop import-all-tables slow和sequence files是自定义java对象
【发布时间】:2017-04-25 17:42:33
【问题描述】:

我正在努力将一个非常大的数据库同步到 hive。

有 2 个问题:(1) 文本导入速度较慢,并且大型 mapreduce 步骤较慢。 (2) 序列文件要快得多,但不能通过正常方式读取。

详情如下:

(1) 如果我们将数据作为文本导入,它会更慢。文件累积在主目录的临时文件夹中,但最终会创建一个相当慢的 mapreduce 作业。

17/04/25 04:18:34 INFO mapreduce.Job: Job job_1490822567992_0996 running in uber mode : false
17/04/25 04:18:34 INFO mapreduce.Job:  map 0% reduce 0%
17/04/25 11:05:59 INFO mapreduce.Job:  map 29% reduce 0%
17/04/25 11:20:18 INFO mapreduce.Job:  map 86% reduce 0% <-- tends to hang a very long time here

(为简洁起见,删除了很多行。)

(2) 如果我们将文件作为序列文件导入,速度会快得多,但 Hive 无法读取检索到的数据,因为它需要了解自动生成的 Java 文件。这也有一个 mapreduce 步骤,但它似乎走得更快(或者也许那是一天中的事情......)。

对于 sqoop 生成的每个表,我们都有一系列这些类: public class MyTableName 扩展 SqoopRecord 实现 DBWritable、Writable

使用这些类的步骤是什么?我们如何在蜂巢中安装它们?令人惊讶的是,Cloudera 支持工程师不知道,因为这一定是不经常绘制的领域?

sqoop import-all-tables --connect '...' --relaxed-isolation --num-mappers 7 --compress --autoreset-to-one-mapper --compression-codec=snappy --outdir javadir --as-sequencefile --hive-delims-replacement ' '

有什么建议吗?

【问题讨论】:

  • "sequencefiles ... Hive 无法读取检索到的数据,因为它需要了解创建的自动生成的 Java 文件" >> 这是什么废话? Hive 只需要一个适当的CREATE TABLE 命令来了解 SequenceFile 结构。这就是--hive-import 的目的。 stackoverflow.com/questions/31515498/…
  • 另外,由于 Sqoop 生成一个 MapReduce 作业,使用 Snappy (或 LZ4) 压缩中间文件和 Snappy (或 LZ4,或 GZip)最终文件的压缩可能会对性能产生重大影响;参看。 mapreduce.map.output.compress*mapreduce.output.fileoutputformat.compress*
  • Hive 导入命令不适用于序列文件。
  • 您是否考虑过使用 Spark 脚本来替代 Sqoop(毕竟,在 Spark 之前的时代,这只是 Cloudera 赞助的权宜之计)我>?作为奖励,您可以获得压缩的 Parquet 文件作为输出。

标签: java performance hadoop sqoop


【解决方案1】:

我对 Spark 持开放态度。你有一些示例代码吗?

免责声明:我刚刚从多个笔记本中组装了一些 sn-ps,并且太懒(也太饿了)无法在离开办公室之前启动测试运行。任何错误和错别字都可以找到。


使用 Cloudera parcel 提供的 Spark 2.0(支持 Hive),这是一种交互式 Scala 脚本,在本地模式下,没有任何数据分区,Microsoft SQL Server 连接,并直接插入到现有的 Hive 托管表(带有一些额外的业务逻辑)......
spark2-shell --master local --driver-class-path /some/path/to/sqljdbc42.jar

// 旁注:4 类 JDBC 驱动程序的自动注册在多个 Spark 构建中被破坏,并且该错误不断再次出现,因此指定驱动程序类以防万一会更安全...

val weather = spark.read.format("jdbc").option("driver", "com.microsoft.sqlserver.jdbc.SQLServerDriver").option("url", "jdbc:sqlserver://myhost\\SQLExpress:9433;database=mydb").option("user", "mylogin").option("password", "*****").option("dbtable", "weather_obs").load()
{ printf( "%%% Partitions: %d / Records: %d\n", weather.rdd.getNumPartitions, weather.count)
  println("%%% Detailed DF schema:")
  weather.printSchema
}

// 使用子查询替代"dbtable" :
//"(SELECT station, dt_obs_utc, temp_k FROM observation_meteo WHERE station LIKE '78%') x")

weather.registerTempTable("wth")
spark.sql(
    """
    INSERT INTO TABLE somedb.sometable
    SELECT station, dt_obs_utc, CAST(temp_k -273.15 AS DECIMAL(3,1)) as temp_c
    FROM wth
    WHERE temp_k IS NOT NULL
    """)
dropTempTable("wth")

weather.unpersist()


现在,如果您想使用 GZip 压缩在 Parquet 文件上动态创建 Hive 外部表,请将“临时表”技巧替换为...
weather.write.option("compression","gzip").mode("overwrite").parquet("hdfs:///some/directory/")

// Parquet 支持的压缩编解码器:无、snappy(默认)、gzip
// 支持的 CSV 压缩编解码器:无(默认)、snappy、lz4、gzip、bzip2

def toImpalaType(sparkType : String ) : String = {
  if (sparkType == "StringType" || sparkType == "BinaryType")  { return "string" }
  if (sparkType == "BooleanType")                              { return "boolean" }
  if (sparkType == "ByteType")                                 { return "tinyint" }
  if (sparkType == "ShortType")                                { return "smallint" }
  if (sparkType == "IntegerType")                              { return "int" }
  if (sparkType == "LongType")                                 { return "bigint" }
  if (sparkType == "FloatType")                                { return "float" }
  if (sparkType == "DoubleType")                               { return "double" }
  if (sparkType.startsWith("DecimalType"))                     { return sparkType.replace("DecimalType","decimal") }
  if (sparkType == "TimestampType" || sparkType == "DateType") { return "timestamp" }
  println("########## ERROR - \"" +sparkType +"\" not supported (bug)")
  return "string"
}

spark.sql("DROP TABLE IF EXISTS somedb.sometable")
{ val query = new StringBuilder
  query.append("CREATE EXTERNAL TABLE somedb.sometable")
  val weatherSchema =weather.dtypes
  val (colName0,colType0) = weatherSchema(0)
  query.append("\n ( " +colName0 + " " +toImpalaType(colType0))
  for ( i <- 2 to tempSchema.length) { val (colName_,colType_) = tempSchema(i-1) ; query.append("\n , " +colName_ + " " +toImpalaType(colType_)) }
  query.append("\n )\nCOMMENT 'Imported from SQL Server by Spark'")
  query.append("\nSTORED AS Parquet")
  query.append("\nLOCATION 'hdfs:///some/directory'")
  sqlContext.sql(query.toString())
  query.clear()
}


如果要对输入表进行分区(基于数字列 - AFAIK 不支持日期/时间),请查看 JDBC 导入选项partitionColumnlowerBoundupperBound

如果您想在 YARN 客户端模式下并行加载这些分区,则添加 --jars 参数以将 JDBC 驱动程序上传到执行程序。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多