相关疑难解决方法(0)

如何在Spark SQL中找到分组Vector列的平均值?

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

5
推荐指数
1
解决办法
2327
查看次数

关于数据集中的 kryo 和 java 编码器的问题

我正在使用 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

3
推荐指数
1
解决办法
2696
查看次数