【发布时间】:2016-09-27 20:03:22
【问题描述】:
我有一个如下所示的 Spark DataFrame:
root
|-- employeeName: string (nullable = true)
|-- employeeId: string (nullable = true)
|-- employeeEmail: string (nullable = true)
|-- company: struct (nullable = true)
| |-- companyName: string (nullable = true)
| |-- companyId: string (nullable = true)
| |-- details: struct (nullable = true)
| | |-- founded: string (nullable = true)
| | |-- address: string (nullable = true)
| | |-- industry: string (nullable = true)
我想做的是按 companyId 分组并为每个公司获取一组员工,如下所示:
root
|-- company: struct (nullable = true)
| |-- companyName: string (nullable = true)
| |-- companyId: string (nullable = true)
| |-- details: struct (nullable = true)
| | |-- founded: string (nullable = true)
| | |-- address: string (nullable = true)
| | |-- industry: string (nullable = true)
|-- employees: array (nullable = true)
| |-- employee: struct (nullable = true)
| | |-- employeeName: string (nullable = true)
| | |-- employeeId: string (nullable = true)
| | |-- employeeEmail: string (nullable = true)
当然,如果我有一对 (company, employee): (String, String) 使用 map 和 reduceByKey,我可以轻松做到这一点。但是对于所有不同的嵌套信息,我不确定要采取什么方法。
我应该尝试压平所有东西吗?任何做类似事情的例子都会非常有帮助。
【问题讨论】:
标签: scala apache-spark