【发布时间】:2019-09-28 16:42:27
【问题描述】:
这个问题本质上是 this question 的重复,除了我在 R 中工作。pyspark 解决方案看起来很可靠,但我无法弄清楚如何将 collect_list 应用于相同的窗口函数在 sparklyr 中的方式。
我有一个具有以下结构的 Spark DataFrame:
------------------------------
userid | date | city
------------------------------
1 | 2018-08-02 | A
1 | 2018-08-03 | B
1 | 2018-08-04 | C
2 | 2018-08-17 | G
2 | 2018-08-20 | E
2 | 2018-08-23 | F
我正在尝试按userid 对DataFrame 进行分组,按date 对每个组进行排序,并将city 列折叠成其值的串联。期望的输出:
------------------
userid | cities
------------------
1 | A, B, C
2 | G, E, F
问题在于,我尝试使用的每种方法都导致一些用户(在对 5000 名用户的测试中约为 3%)的“城市”列的顺序不正确。
尝试 1:使用 dplyr 和 collect_list。
my_sdf %>%
dplyr::group_by(userid) %>%
dplyr::arrange(date) %>%
dplyr::summarise(cities = paste(collect_list(city), sep = ", ")))
尝试2:使用replyr::gapply,因为该操作符合“Grouped-Order-Apply”的描述。
get_cities <- . %>%
summarise(cities = paste(collect_list(city), sep = ", "))
my_sdf %>%
replyr::gapply(gcolumn = "userid",
f = get_cities,
ocolumn = "date",
partitionMethod = "group_by")
尝试 3:编写为 SQL 窗口函数。
my_sdf %>%
spark_session(sc) %>%
sparklyr::invoke("sql",
"SELECT userid, CONCAT_WS(', ', collect_list(city)) AS cities
OVER (PARTITION BY userid
ORDER BY date)
FROM my_sdf") %>%
sparklyr::sdf_register() %>%
sparklyr::sdf_copy_to(sc, ., "my_sdf", overwrite = T)
^ 抛出以下错误:
Error: org.apache.spark.sql.catalyst.parser.ParseException:
mismatched input 'OVER' expecting <EOF>(line 2, pos 19)
== SQL ==
SELECT userid, conversion_location, CONCAT_WS(' > ', collect_list(channel)) AS path
OVER (PARTITION BY userid, conversion_location
-------------------^^^
ORDER BY occurred_at)
FROM paths_model
【问题讨论】: