【问题标题】:What is the most efficient way to create new Spark Tables or Data Frames in Sparklyr?在 Sparklyr 中创建新的 Spark 表或数据框的最有效方法是什么?
【发布时间】:2017-11-27 08:03:43
【问题描述】:

在 Hadoop 集群(不是虚拟机)上使用 sparklyr 包,我正在处理需要连接、过滤等的几种类型的表......我正在尝试确定什么是使用dplyr 命令以及sparklyr 中的数据管理功能来运行处理、将其存储在缓存中并与中间数据对象一起使用以生成保留在缓存中的下游对象的最有效方法。这个问题就像上面提出的那样肤浅,但我希望获得比纯粹效率更多的信息,所以如果你想编辑我的问题,我可以接受......

我在 Hive 中有一些表,我们称它们为 Activity2016Accounts2016Accounts2017。 “帐户”表还包括地址历史记录。我想从 2016 年的数据开始,合并姓名和当前地址的两个表,过滤一些活动和帐户详细信息,然后将两种不同的方式与 2017 年的帐户信息合并,特别是统计留在他们地址的人数与那些谁换了地址。我们有数百万行,因此我们正在使用我们的 spark 集群进行此活动。

所以,首先,这就是我现在正在做的事情:

sc <- spark_connect()    

Activity2016 %>% filter(COL1 < Cut1 & COL1 > Cut2) %>% 
select(NAME,ADDRESS1) %>% 
inner_join(Accounts2016,c("NAME"="NAME","ADDRESS1"="ADDRESS1")) %>%
distinct(NAME,ADDRESS1) %>% sdf_register("JOIN2016")

tbl_cache(sc,"JOIN2016")
JOINED_2016 <- tbl(sc, "JOIN2016")

Acct2017 = tbl(sc, "HiveDB.Accounts2017")

# Now, I can run:
JOINED_2016 %>% inner_join(Acct2017,c("NAME"="NAME","ADDRESS1"="ADDRESS2")) %>%
distinct(NAME,ADDRESS1.x) %>% sdf_register("JOIN2017")

# Rinse & Repeat
tbl_cache(sc,"JOIN2017")
JOINED_2017 <- tbl(sc,"JOIN2017")

然后我继续使用JOINED_2016JOINED_2017,使用dplyr 动词等...

这里似乎有多个效率低下的地方......比如, 1)我不应该能够将它直接发送到缓存并将其作为变量调用吗? 2)我不应该也可以将它直接发送到书面的 Hive 表中吗? 3) 如何转换最终对象以运行基本的R 命令,例如table(JOINED_2016$COL1),或者这些命令不可用(尝试%&gt;% select(COL1) %&gt;% table 时出现错误)?

如果有下游错误我会丢失数据并且我不写它......但是我觉得关于如何写我不清楚的数据的选择太多了。它何时最终成为缓存对象,与RDD,与 Hive 内部/外部表,与 Spark DataFrame 相比,R 处理这些数据对象的能力有哪些限制?

例如,如果我只是运行:

JOIN2016 <- Activity2016 %>% filter(COL1 < Cut1 & COL1 > Cut2) %>% 
select(NAME,ADDRESS1) %>% 
inner_join(Accounts2016,c("NAME"="NAME","ADDRESS1"="ADDRESS1")) %>%
distinct(NAME,ADDRESS1) 

这是一个 R data.frame 吗? (这可能会导致我的网关节点的 RAM 崩溃......这就是为什么我不愿意尝试它。这是企业中的一个集群)

总结一下: 我应该完全使用tbltbl_cache 命令还是需要它们?

我应该使用dbWriteTable,我可以在之后、之前或代替sdf_register 直接使用它吗?还是我需要使用tbl 命令才能向Hive 写入任何内容? sdf_register 几乎毫无意义。

我应该使用copy_to 还是db_copy_to 而不是dbWriteTable?我不想把 Hive 变成垃圾场,所以我要小心我如何编写中间数据,然后在我存储它之后保持 R 将如何使用它。

我必须运行哪些data.frame-type 来处理数据,就像它是内存中的 R 对象一样,还是我仅限于 dplyr 命令?

很抱歉,这个问题的内容太多了,但我觉得 R-bloggers 文章、sparklyr 教程以及 SOF 上的其他问题都不清楚这些问题。

【问题讨论】:

    标签: hadoop apache-spark hive dplyr sparklyr


    【解决方案1】:

    sdf_register 在处理长时间运行的查询时不是很有用。它基本上是一个非物化视图,这意味着每次调用它时它都会运行底层查询。添加以下内容会将数据作为表写入 Hive。

    spark_dataframe %&gt;% invoke("write") %&gt;% invoke("saveAsTable", as.character("your_desired_table_name"))

    这使用saveAsTable 作为表,这将创建一个表,即使在 Spark 会话结束后也会保留该表。当 Spark 会话结束时,使用 createOrReplaceTempView 不会保留数据。

    【讨论】:

    • 这是否意味着,在断开连接后重新连接时,需要运行一个新呼叫,即desired_tbl &lt;- tbl("your_desired_table_name")
    • 是的。从这个过程中休息的表就像所有其他永久 Hive 表一样。
    • 另外,当我检查这个对象的class 时,我得到"tbl_spark" "tbl_sql" "tbl_lazy" "tbl"。为了调用适用于data.frame 的函数,看起来我们唯一能做的就是运行mydf &lt;- collect(desired_tbl)
    猜你喜欢
    • 1970-01-01
    • 2022-08-13
    • 1970-01-01
    • 2012-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-05-29
    相关资源
    最近更新 更多