【问题标题】:Multiple SQLs on a temp table failed临时表上的多个 SQL 失败
【发布时间】:2017-06-30 02:53:58
【问题描述】:
Spark Version: 1.6.2.   

我注册了一个临时表,数据源是HDFS,查询了两次。

然后作业因以下错误而失败:

错误 ApplicationMaster:用户类抛出异常:
java.io.IOException:不是文件:hdfs://my_server:8020/2017/01/01
java.io.IOException:不是文件:hdfs://my_server:8020/2017/01/01 在 org.apache.hadoop.mapred.FileInputFormat.getSplits(FileInputFormat.java:322) 在 org.apache.spark.rdd.HadoopRDD.getPartitions(HadoopRDD.scala:199) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:242) 在 org.apache.spark.rdd.RDD$$anonfun$partitions$2.apply(RDD.scala:240) 在 scala.Option.getOrElse(Option.scala:120) 在 org.apache.spark.rdd.RDD.partitions(RDD.scala:240)

棘手的部分是,如果只运行一个查询,作业就会成功。
我是否以错误的方式使用 Spark SQL,或者这是故意的?

这是我的代码的样子:

val rdd = sc.textFile("hdfs://my_server:8020/2017/*/*/*")
val table = sqlc.read.json(rdd).cache()

table.registerTempTable("my_table")

sql("""
    | SELECT contentsId,
    |   SUM(CASE WHEN gender = 'M' then 1 else 0 end)
    | FROM my_table
    | GROUP BY contentsId
  """.stripMargin)
  .write.format("com.databricks.spark.csv")
  .save("hdfs://my_server:8020/gender.csv")

sql("""
    | SELECT contentsId,
    |   SUM(CASE WHEN age > 0 AND age < 20 then 1 else 0 end),
    |   SUM(CASE WHEN age >= 20 AND age < 30 then 1 else 0 end)
    | FROM my_table
    | GROUP BY contentsId
  """.stripMargin)
  .write.format("com.databricks.spark.csv")
  .save("hdfs://my_server:8020/age.csv")

提前致谢!

【问题讨论】:

  • 您为什么要在一个输出文件中将两个不同的数据帧保存为 gender.csv?
  • 错误表示无法读取文件 不是文件:hdfs://my_server:8020/2017/01/01
  • @ShankarKoirala 我的错。我需要将它们保存在两个不同的文件中。我编辑了查询。
  • @Bhavesh hdfs://my_server:8020/2017/01/01 是一个目录。而且,这个异常只发生在一个作业中有两个查询时。

标签: apache-spark apache-spark-sql spark-dataframe


【解决方案1】:

我认为您可以尝试仅对此类文件应用过滤器。

val filesRDD = rdd.filter{path => (new java.io.File(path).isFile)}

这将删除 RDD 中包含的所有目录 并且第二次保存 DataFrame 使用这个

sql("""
    | SELECT contentsId,
    |   SUM(CASE WHEN age > 0 AND age < 20 then 1 else 0 end),
    |   SUM(CASE WHEN age >= 20 AND age < 30 then 1 else 0 end)
    | FROM my_table
    | GROUP BY contentsId
  """.stripMargin)
  .write.format("com.databricks.spark.csv")
  .mode("append")
  .save("hdfs://my_server:8020/gender.csv")

如果存储值相同或尝试将第二个 DataFrame 存储到其他文件中

【讨论】:

  • 过滤路径是个好主意,但关键是只有当两个查询在一个工作中时才会发生异常。 :-( 你能添加更多省时部分的声明吗?
猜你喜欢
  • 1970-01-01
  • 2011-02-08
  • 1970-01-01
  • 2018-04-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-07-21
  • 2010-12-21
相关资源
最近更新 更多