我RelationalGroupedDataset通过调用创建了一个instances.groupBy(instances.col("property_name")):
val x = instances.groupBy(instances.col("property_name"))
Run Code Online (Sandbox Code Playgroud)
如何组合用户定义的聚合函数来对每个组执行Statistics.colStats().mean?
谢谢!
aggregate-functions user-defined-functions apache-spark apache-spark-sql apache-spark-ml
我正在使用 Spark 2.4 并参考 https://spark.apache.org/docs/latest/rdd-programming-guide.html#rdd-persistence
豆类:
public class EmployeeBean implements Serializable {
private Long id;
private String name;
private Long salary;
private Integer age;
// getters and setters
}
Run Code Online (Sandbox Code Playgroud)
火花示例:
SparkSession spark = SparkSession.builder().master("local[4]").appName("play-with-spark").getOrCreate();
List<EmployeeBean> employees1 = populateEmployees(1, 1_000_000);
Dataset<EmployeeBean> ds1 = spark.createDataset(employees1, Encoders.kryo(EmployeeBean.class));
Dataset<EmployeeBean> ds2 = spark.createDataset(employees1, Encoders.bean(EmployeeBean.class));
ds1.persist(StorageLevel.MEMORY_ONLY());
long ds1Count = ds1.count();
ds2.persist(StorageLevel.MEMORY_ONLY());
long ds2Count = ds2.count();
Run Code Online (Sandbox Code Playgroud)
我在 Spark Web UI 中寻找存储。有用的部分——
ID RDD Name Size in Memory
2 LocalTableScan [value#0] 56.5 MB
13 LocalTableScan [age#6, id#7L, name#8, …Run Code Online (Sandbox Code Playgroud) kryo apache-spark apache-spark-dataset apache-spark-encoders