【发布时间】:2015-10-27 18:48:39
【问题描述】:
我正在尝试使用 SparkSql HiveContext 读取 Hive 表。但是,当我提交作业时,我收到以下错误:
Exception in thread "main" java.lang.RuntimeException: Unsupported parquet datatype optional fixed_len_byte_array(11) amount (DECIMAL(24,7))
at scala.sys.package$.error(package.scala:27)
at org.apache.spark.sql.parquet.ParquetTypesConverter$.toPrimitiveDataType(ParquetTypes.scala:77)
at org.apache.spark.sql.parquet.ParquetTypesConverter$.toDataType(ParquetTypes.scala:131)
at org.apache.spark.sql.parquet.ParquetTypesConverter$$anonfun$convertToAttributes$1.apply(ParquetTypes.scala:383)
at org.apache.spark.sql.parquet.ParquetTypesConverter$$anonfun$convertToAttributes$1.apply(ParquetTypes.scala:380)
列类型为 DECIMAL(24,7)。我已经使用 HiveQL 更改了列类型,但它不起作用。我还尝试在 sparksql 中转换为另一种 Decimal 类型,如下所示:
val results = hiveContext.sql("SELECT cast(amount as DECIMAL(18,7)), number FROM dmp_wr.test")
但是,我得到了同样的错误。我的代码是这样的:
def main(args: Array[String]) {
val conf: SparkConf = new SparkConf().setAppName("TColumnModify")
val sc: SparkContext = new SparkContext(conf)
val vectorAcc = sc.accumulator(new MyVector())(VectorAccumulator)
val hiveContext = new org.apache.spark.sql.hive.HiveContext(sc)
val results = hiveContext.sql("SELECT amount, number FROM dmp_wr.test")
我该如何解决这个问题?感谢您的回复。
Edit1:我找到了引发异常的 Spark 源代码行。好像是这样的
if(originalType == ParquetOriginalType.DECIMAL && decimalInfo.getPrecision <= 18)
因此,我创建了一个新表,其中包含 DECIMAL(18,7) 类型的列,并且我的代码按预期工作。
我删除表并创建具有 DECIMAL(24,7) 列的新表,之后我更改了列类型
alter table qwe change amount amount decimal(18,7) 我可以看到它已更改为 DECIMAL(18,7),但 Spark
不接受改变。它仍然将列类型读取为 DECIMAL(24,7) 并给出相同的错误。
主要原因是什么?
【问题讨论】:
标签: scala apache-spark hive cloudera apache-spark-sql