dhe*_*eee 10 scala apache-spark apache-spark-sql spark-dataframe
我正在尝试计算DataFrame中列的百分位数?我无法在Spark聚合函数中找到任何percentile_approx函数.
例如在Hive中我们有percentile_approx,我们可以通过以下方式使用它
hiveContext.sql("select percentile_approx("Open_Rate",0.10) from myTable);
Run Code Online (Sandbox Code Playgroud)
但出于性能原因,我想使用Spark DataFrame来实现它.
样本数据集
|User ID|Open_Rate|
-------------------
|A1 |10.3 |
|B1 |4.04 |
|C1 |21.7 |
|D1 |18.6 |
Run Code Online (Sandbox Code Playgroud)
我想知道有多少用户分为10百分位或20百分位等等.我想做这样的事情
df.select($"id",Percentile($"Open_Rate",0.1)).show
Run Code Online (Sandbox Code Playgroud)
从Spark2.0开始,事情变得越来越容易,只需在DataFrameStatFunctions中使用此函数,例如:
df.stat.approxQuantile("Open_Rate",Array(0.25,0.50,0.75),0.0)
DataFrameStatFunctions中的DataFrame还有一些有用的统计函数.
SparkSQL 和 Scala 数据帧/数据集 API 由同一引擎执行。等效的操作将生成等效的执行计划。您可以使用 来查看执行计划explain。
sql(...).explain
df.explain
Run Code Online (Sandbox Code Playgroud)
当谈到您的具体问题时,混合 SparkSQL 和 Scala DSL 语法是一种常见模式,因为正如您所发现的,它们的功能尚未等效。(另一个例子是 SQLexplode()和 DSL之间的区别explode(),后者更强大,但由于编组而效率更低。)
简单的方法如下:
df.registerTempTable("tmp_tbl")
val newDF = sql(/* do something with tmp_tbl */)
// Continue using newDF with Scala DSL
Run Code Online (Sandbox Code Playgroud)
如果您采用简单的方法,需要记住的是临时表名称是集群全局的(最高 1.6.x)。因此,如果代码可能在同一集群上同时运行多次,则应使用随机表名称。
在我的团队中,这种模式很常见,以至于我们添加了一个.sql()隐式函数DataFrame,可以在 SQL 语句的范围内自动注册然后取消注册临时表。
| 归档时间: |
|
| 查看次数: |
9831 次 |
| 最近记录: |