在Spark SQL中使用collect_list和collect_set

Joo*_*rla 15 hive apache-spark apache-spark-sql

根据文档,这些collect_set和collect_list函数应该在Spark SQL中可用.但是,我无法让它发挥作用.我正在使用Docker镜像运行Spark 1.6.0 .

我想在Scala中这样做:

import org.apache.spark.sql.functions._ 

df.groupBy("column1") 
  .agg(collect_set("column2")) 
  .show() 
Run Code Online (Sandbox Code Playgroud)

并在运行时收到以下错误:

Exception in thread "main" org.apache.spark.sql.AnalysisException: undefined function collect_set; 
Run Code Online (Sandbox Code Playgroud)

也尝试使用它pyspark,但它也失败了.文档声明这些函数是Hive UDAF的别名,但我无法想出启用这些函数.

如何解决这个问题?感谢名单!

zer*_*323 33

Spark 2.0+:

SPARK-10605引入了原生collect_list和collect_set实现.SparkSession有Hive支持或HiveContext不再需要.

Spark 2.0-SNAPSHOT (2016-05-03之前):

您必须为给定的Hive支持启用SparkSession:

在斯卡拉:

val spark = SparkSession.builder
  .master("local")
  .appName("testing")
  .enableHiveSupport()  // <- enable Hive support.
  .getOrCreate()
Run Code Online (Sandbox Code Playgroud)

在Python中:

spark = (SparkSession.builder
    .enableHiveSupport()
    .getOrCreate())
Run Code Online (Sandbox Code Playgroud)

Spark <2.0:

为了能够使用Hive UDF(请参阅https://cwiki.apache.org/confluence/display/Hive/LanguageManual+UDF),您已经使用了使用Hive支持构建的Spark(当您使用预先构建的二进制文件时,这已经涵盖了似乎是这里的情况)并初始化SparkContext使用HiveContext.

在斯卡拉:

import org.apache.spark.sql.hive.HiveContext
import org.apache.spark.sql.SQLContext

val sqlContext: SQLContext = new HiveContext(sc) 
Run Code Online (Sandbox Code Playgroud)

在Python中:

from pyspark.sql import HiveContext

sqlContext = HiveContext(sc)
Run Code Online (Sandbox Code Playgroud)