【发布时间】:2020-05-07 01:14:30
【问题描述】:
我正在寻找有关如何解决以下情况的想法。我的用例是在 java spark 中,但是在我用尽想法时寻找有关如何做到这一点的想法,而不管语言如何
我有如下非结构化数据
98480|PERSON|TOM|GREER|1982|12|27
98480|PHONE|CELL|732|201|6789
98480|PHONE|HOME|732|123|9876
98480|ADDR|RES|102|JFK BLVD|PISCATAWAY|NJ|08854
98480|ADDR|OFF|211|EXCHANGE PL|JERSEY CITY|NJ|07302
98481|PERSON|LIN|JASSOY|1976|09|15
98481|PHONE|CELL|908|398|3389
98481|PHONE|HOME|917|363|2647
98481|ADDR|RES|111|JOURNAL SQ|JERSEY CITY|NJ|07704
98481|ADDR|OFF|365|DOWNTOWN NEWYORK|NEWYORK CITY|NY|10001
我正在尝试将它们转换为带有 persondata 的行,其中包含一组电话和 addr,如下所示,每个 personId 基本上是单行
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------------+
|personId|type |firstName|lastName|year|month|day|Phone | addr | |
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------------+
|98481 |PERSON|LIN |JASSOY |1976|09 |15 |[[PHONE, HOME, 917, 363, 2647], [PHONE, CELL, 908, 398, 3389]] | [[ADDR, OFF, 365, DOWNTOWN NEWYORK, NEWYORK CITY, NY, 10001], [ADDR, RES, 111, JOURNAL SQ, JERSEY CITY, NJ, 07704]] |
|98480 |PERSON|TOM |GREER |1982|12 |27 |[[PHONE, HOME, 732, 123, 9876], [PHONE, CELL, 732, 201, 6789]] | [[ADDR, RES, 102, JFK BLVD, PISCATAWAY, NJ, 08854], [ADDR, OFF, 211, EXCHANGE PL, JERSEY CITY, NJ, 07302]] |
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------------+-----------------------------------------------------------------------------------------------------------------------+
下面的代码
Dataset<Row> dataset = groupedDataset
.agg(collect_set(struct(phoneRow.col("type").as("collType"), phoneRow.col("phoneType").as("phoneType"),
phoneRow.col("areaCode").as("areaCode"), phoneRow.col("phoneMiddle").as("phoneMiddle"),
phoneRow.col("ext").as("ext"), addressRow.col("type").as("collType"),
addressRow.col("addrType").as("addrType"), addressRow.col("addr1").as("rowType"),
addressRow.col("addr2").as("addr2"), addressRow.col("city").as("city"),
addressRow.col("state").as("state"), addressRow.col("zipCode").as("zipCode"))).as("addrPhone"));
输出如下,但不是我要的格式
+--------+------+---------+--------+----+-----+---+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
|personId|type |firstName|lastName|year|month|day|addrPhone |
+--------+------+---------+--------+----+-----+---+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
|98481 |PERSON|LIN |JASSOY |1976|09 |15 |[[PHONE, HOME, 917, 363, 2647, ADDR, OFF, 365, DOWNTOWN NEWYORK, NEWYORK CITY, NY, 10001], [PHONE, HOME, 917, 363, 2647, ADDR, RES, 111, JOURNAL SQ, JERSEY CITY, NJ, 07704], [PHONE, CELL, 908, 398, 3389, ADDR, RES, 111, JOURNAL SQ, JERSEY CITY, NJ, 07704], [PHONE, CELL, 908, 398, 3389, ADDR, OFF, 365, DOWNTOWN NEWYORK, NEWYORK CITY, NY, 10001]]|
|98480 |PERSON|TOM |GREER |1982|12 |27 |[[PHONE, HOME, 732, 123, 9876, ADDR, RES, 102, JFK BLVD, PISCATAWAY, NJ, 08854], [PHONE, CELL, 732, 201, 6789, ADDR, RES, 102, JFK BLVD, PISCATAWAY, NJ, 08854], [PHONE, CELL, 732, 201, 6789, ADDR, OFF, 211, EXCHANGE PL, JERSEY CITY, NJ, 07302], [PHONE, HOME, 732, 123, 9876, ADDR, OFF, 211, EXCHANGE PL, JERSEY CITY, NJ, 07302]] |
+--------+------+---------+--------+----+-----+---+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
寻找解决上述问题的想法
更新: 我能够按预期获得输出,但我不确定它的效果如何,并且看起来有很多带有大量连接和数据框的样板代码。这是我用来理解 spark 的示例数据,但我要处理的真实数据会有很多复杂的转换,而且这段代码看起来并不有效
这里是更新的代码
Dataset<Row> groupedPhoneDataSet = groupedDataset.agg(collect_set(struct(phoneRow.col("type").as("phColType"),
phoneRow.col("phoneType").as("phoneType"), phoneRow.col("areaCode").as("areaCode"),
phoneRow.col("phoneMiddle").as("phoneMiddle"), phoneRow.col("ext").as("ext"))).as("phoneRec"));
Dataset<Row> groupedAddrDataSet = groupedDataset
.agg(collect_set(struct(addressRow.col("type").as("addrColType"),
addressRow.col("addrType").as("addrType"), addressRow.col("addr1").as("addr1"),
addressRow.col("addr2").as("addr2"), addressRow.col("city").as("city"),
addressRow.col("state").as("state"), addressRow.col("zipCode").as("zipCode"))).as("addrRec"));
Dataset<Row> finalDataSet = groupedAddrDataSet
.join(groupedPhoneDataSet,
groupedAddrDataSet.col("personId").equalTo(groupedPhoneDataSet.col("personId")))
.select(groupedPhoneDataSet.col("personId"), groupedPhoneDataSet.col("type"),
groupedPhoneDataSet.col("firstName"), groupedPhoneDataSet.col("lastName"),
groupedPhoneDataSet.col("year"), groupedPhoneDataSet.col("month"),
groupedPhoneDataSet.col("day"), col("phoneRec"), col("addrRec"));
这是我得到的输出
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------+
|personId|type |firstName|lastName|year|month|day|phoneRec |addrRec |
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------+
|98481 |PERSON|LIN |JASSOY |1976|09 |15 |[[PHONE, CELL, 908, 398, 3389], [PHONE, HOME, 917, 363, 2647]]|[[ADDR, RES, 111, JOURNAL SQ, JERSEY CITY, NJ, 07704], [ADDR, OFF, 365, DOWNTOWN NEWYORK, NEWYORK CITY, NY, 10001]]|
|98480 |PERSON|TOM |GREER |1982|12 |27 |[[PHONE, CELL, 732, 201, 6789], [PHONE, HOME, 732, 123, 9876]]|[[ADDR, OFF, 211, EXCHANGE PL, JERSEY CITY, NJ, 07302], [ADDR, RES, 102, JFK BLVD, PISCATAWAY, NJ, 08854]] |
+--------+------+---------+--------+----+-----+---+--------------------------------------------------------------+-------------------------------------------------------------------------------------------------------------------+
有没有办法在不创建大量数据帧的情况下做到这一点
【问题讨论】:
标签: java apache-spark apache-spark-sql