【问题标题】:Process big data using hadoop parquet to CSV output使用 hadoop parquet 处理大数据到 CSV 输出
【发布时间】:2017-10-16 16:01:42
【问题描述】:

我有 3 个数据集,我想将它们加入并分组,以获得包含聚合数据的 CSV。

数据作为 parquet 文件存储在 Hadoop 中,我使用 Zeppelin 运行 Apache Spark+Scala 进行数据处理。

我的数据集如下所示:

user_actions.show(10)
user_clicks.show(10)
user_options.show(10)

+--------------------+--------------------+
|                  id|             keyword|
+--------------------+--------------------+
|00000000000000000001|               aaaa1|
|00000000000000000002|               aaaa1|
|00000000000000000003|               aaaa2|
|00000000000000000004|               aaaa2|
|00000000000000000005|               aaaa0|
|00000000000000000006|               aaaa4|
|00000000000000000007|               aaaa1|
|00000000000000000008|               aaaa2|
|00000000000000000009|               aaaa1|
|00000000000000000010|               aaaa1|
+--------------------+--------------------+
+--------------------+-------------------+
|           search_id|   selected_user_id|
+--------------------+-------------------+
|00000000000000000001|               1234|
|00000000000000000002|               1234|
|00000000000000000003|               1234|
|00000000000000000004|               1234|
+--------------------+-------------------+

+--------------------+----------+----------+
|           search_id|   user_id|  position|
+--------------------+----------+----------+
|00000000000000000001|      1230|         1|
|00000000000000000001|      1234|         3|
|00000000000000000001|      1232|         2|
|00000000000000000002|      1231|         1|
|00000000000000000002|      1232|         2|
|00000000000000000002|      1233|         3|
|00000000000000000002|      1234|         4|
|00000000000000000003|      1234|         1|
|00000000000000000004|      1230|         1|
|00000000000000000004|      1234|         2|
+--------------------+----------+----------+

我想要实现的是为每个用户 ID 获取一个带有关键字的 JSON,因为我需要将它们导入 MySQL 并将 user_id 作为 PK。

user_id,keywords
1234,"{\"aaaa1\":3.5,\"aaaa2\":0.5}"

如果 JSON 不是开箱即用的,我可以使用元组或任何字符串:

user_id,keywords
1234,"(aaaa1,0.58333),(aaaa2,1.5)"

到目前为止我所做的是:

val user_actions_data = user_actions
                                .join(user_options, user_options("search_id") === user_actions("id"))

val user_actions_full_data = user_actions_data
                                    .join(
                                            user_clicks,
                                            user_clicks("search_id") === user_actions_data("search_id") && user_clicks("selected_user_id") === user_actions_data("user_id"),
                                            "left_outer"
                                        )

val user_actions_data_groupped = user_actions_full_data
                                        .groupBy("user_id", "search")
                                        .agg("search" -> "count", "selected_user_id" -> "count", "position" -> "avg")


def udfScoreForUser = ((position: Double, searches: Long) =>  ( position/searches ))

val search_log_keywords = user_actions_data_groupped.rdd.map({row => row(0) -> (row(1) -> udfScoreForUser(row.getDouble(4), row.getLong(2)))}).groupByKey()


val search_log_keywords_array = search_log_keywords.collect.map(r => (r._1.asInstanceOf[Long], r._2.mkString(", ")))

val search_log_keywords_df = sc.parallelize(search_log_keywords_array).toDF("user_id","keywords")
    .coalesce(1)
    .write.format("csv")
    .option("header", "true")
    .mode("overwrite")
    .save("hdfs:///Search_log_testing_keywords/")

虽然这对小型数据集按预期工作,但我的输出 CSV 文件是:

user_id,keywords
1234,"(aaaa1,0.58333), (aaaa2,0.5)"

我在运行 200+GB 数据时遇到了问题。

我对 Spark 和 Scala 还很陌生,但我认为我遗漏了一些东西,我不应该使用 DF 进行 rdd、收集以映射到数组并将其并行化回 DF 以将其导出为 CSV。

总而言之,我想对所有关键字应用评分并按用户 ID 对它们进行分组并将其保存到 CSV。到目前为止,我所做的适用于一个小数据集,但是当我将其应用于 200GB 以上的数据时,apache spark 失败了。

【问题讨论】:

    标签: scala hadoop apache-spark apache-zeppelin


    【解决方案1】:

    是的,任何依赖于 Spark 中 collect 的东西通常都是错误的——除非你正在调试某些东西。当您调用 collect 时,所有数据都在驱动程序中收集到一个数组中,因此对于大多数大数据集,这甚至不是一个选项 - 您的驱动程序将抛出 OOM 并死掉。

    我不明白的是,你为什么要收藏?为什么不简单地映射分布式数据集?

    search_log_keywords
      .map(r => (r._1.asInstanceOf[Long], r._2.mkString(", ")))
      .toDF("user_id","keywords")
      .coalesce(1)
      .write.format("csv")
      .option("header", "true")
      .mode("overwrite")
      .save("hdfs:///Search_log_testing_keywords/")
    

    这样,一切都是并行进行的。

    关于在dataframesrdds 之间切换,我现在不会太担心。我知道社区大多提倡使用dataframes,但根据 Spark 的版本和您的用例,rdds 可能是更好的选择。

    【讨论】:

    • 哦,不,我认为解决了它。我没有看到任何错误,但需要一段时间才能遍历整个数据。你有什么线索我可以得到JSON而不是元组字符串吗?
    • write.json(your/target/path) 怎么样?
    • 当然可以。这确实很有帮助。 collect 把事情搞砸了。但是,我无法完全运行它。我有这个错误:Caused by: java.lang.RuntimeException: java.io.FileNotFoundException: /yarn/nm/usercache/hdfs/appcache/application_1493881537049_0724/blockmgr-933c832c-636a-4748-9103-eff656f2456c/25/shuffle_92_0_0.index (No such file or directory) 我想我运行它的数据仍然比我的服务器可以处理的多。 :(
    • 您使用的是coalesce(1),这意味着所有数据将被放入单个分区并写入单个文件。除非您知道所有数据都适合单个分区,否则您希望避免此类行为。
    • 你得到的错误,通常意味着执行者已经死亡,因此无法找到洗牌数据。执行者通常死于 OOM 异常。如果您仔细检查日志文件,我相信您会发现 OOM 错误。
    【解决方案2】:

    HDFS 的主要目标是将文件分割成块并冗余存储。除非绝对有必要拥有一个大文件,否则最好将数据分区存储在 HDFS 中。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2012-03-04
      • 1970-01-01
      • 2012-07-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2010-11-13
      相关资源
      最近更新 更多