如何通过加入Spark创建嵌套列?

sea*_*avi 1 scala apache-spark apache-spark-sql

我想在两个Spark DataFrames(Scala)上执行“联接”,但是我想将第二个DataFrame的“ joined”行作为单个嵌套列插入第一个,而不是类似SQL的联接。这样做的原因最终是使用嵌套结构写回JSON。我知道答案可能已经在Stackoverflow上了,但是一些搜索并没有找到我的答案。

表格1

root
 |-- Insdc: string (nullable = true)
 |-- LastMetaUpdate: string (nullable = true)
 |-- LastUpdate: string (nullable = true)
 |-- Published: string (nullable = true)
 |-- Received: string (nullable = true)
 |-- ReplacedBy: string (nullable = true)
 |-- Status: string (nullable = true)
 |-- Type: string (nullable = true)
 |-- accession: string (nullable = true)
 |-- alias: string (nullable = true)
 |-- attributes: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- tag: string (nullable = true)
 |    |    |-- value: string (nullable = true)
 |-- center_name: string (nullable = true)
 |-- design_description: string (nullable = true)
 |-- geo_accession: string (nullable = true)
 |-- instrument_model: string (nullable = true)
 |-- library_construction_protocol: string (nullable = true)
 |-- library_name: string (nullable = true)
 |-- library_selection: string (nullable = true)
 |-- library_source: string (nullable = true)
 |-- library_strategy: string (nullable = true)
 |-- paired: boolean (nullable = true)
 |-- platform: string (nullable = true)
 |-- read_spec: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- base_coord: long (nullable = true)
 |    |    |-- read_class: string (nullable = true)
 |    |    |-- read_index: long (nullable = true)
 |    |    |-- read_type: string (nullable = true)
 |-- sample_accession: string (nullable = true)
 |-- spot_length: long (nullable = true)
 |-- study_accession: string (nullable = true)
 |-- tags: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- title: string (nullable = true)
Run Code Online (Sandbox Code Playgroud)

表2

root
 |-- BioProject: string (nullable = true)
 |-- Insdc: string (nullable = true)
 |-- LastMetaUpdate: string (nullable = true)
 |-- LastUpdate: string (nullable = true)
 |-- Published: string (nullable = true)
 |-- Received: string (nullable = true)
 |-- ReplacedBy: string (nullable = true)
 |-- Status: string (nullable = true)
 |-- Type: string (nullable = true)
 |-- abstract: string (nullable = true)
 |-- accession: string (nullable = true)
 |-- alias: string (nullable = true)
 |-- attributes: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- tag: string (nullable = true)
 |    |    |-- value: string (nullable = true)
 |-- dbGaP: string (nullable = true)
 |-- description: string (nullable = true)
 |-- external_id: struct (nullable = true)
 |    |-- id: string (nullable = true)
 |    |-- namespace: string (nullable = true)
 |-- submitter_id: struct (nullable = true)
 |    |-- id: string (nullable = true)
 |    |-- namespace: string (nullable = true)
 |-- tags: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- title: string (nullable = true)
Run Code Online (Sandbox Code Playgroud)

联接table1.study_accession与一起进行table2.accession。结果如下。注意新的列称为study包含表2中行的记录等效项。

root
 |-- Insdc: string (nullable = true)
 |-- LastMetaUpdate: string (nullable = true)
 |-- LastUpdate: string (nullable = true)
 |-- Published: string (nullable = true)
 |-- Received: string (nullable = true)
 |-- ReplacedBy: string (nullable = true)
 |-- Status: string (nullable = true)
 |-- Type: string (nullable = true)
 |-- accession: string (nullable = true)
 |-- alias: string (nullable = true)
 |-- attributes: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- tag: string (nullable = true)
 |    |    |-- value: string (nullable = true)
 |-- center_name: string (nullable = true)
 |-- design_description: string (nullable = true)
 |-- geo_accession: string (nullable = true)
 |-- instrument_model: string (nullable = true)
 |-- library_construction_protocol: string (nullable = true)
 |-- library_name: string (nullable = true)
 |-- library_selection: string (nullable = true)
 |-- library_source: string (nullable = true)
 |-- library_strategy: string (nullable = true)
 |-- paired: boolean (nullable = true)
 |-- platform: string (nullable = true)
 |-- read_spec: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- base_coord: long (nullable = true)
 |    |    |-- read_class: string (nullable = true)
 |    |    |-- read_index: long (nullable = true)
 |    |    |-- read_type: string (nullable = true)
 |-- sample_accession: string (nullable = true)
 |-- spot_length: long (nullable = true)
 |-- study_accession: string (nullable = true)
 |-- tags: array (nullable = true)
 |    |-- element: string (containsNull = true)
 |-- title: string (nullable = true)
 |-- accession: string (nullable = true)
 |-- study: struct (nullable = true)
 |    |-- BioProject: string (nullable = true)
 |    |-- Insdc: string (nullable = true)
 |    |-- LastMetaUpdate: string (nullable = true)
 |    |-- LastUpdate: string (nullable = true)
 |    |-- Published: string (nullable = true)
 |    |-- Received: string (nullable = true)
 |    |-- ReplacedBy: string (nullable = true)
 |    |-- Status: string (nullable = true)
 |    |-- Type: string (nullable = true)
 |    |-- abstract: string (nullable = true)
 |    |-- accession: string (nullable = true)
 |    |-- alias: string (nullable = true)
 |    |-- attributes: array (nullable = true)
 |    |    |-- element: struct (containsNull = true)
 |    |    |    |-- tag: string (nullable = true)
 |    |    |    |-- value: string (nullable = true)
 |    |-- dbGaP: string (nullable = true)
 |    |-- description: string (nullable = true)
 |    |-- external_id: struct (nullable = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- namespace: string (nullable = true)
 |    |-- submitter_id: struct (nullable = true)
 |    |    |-- id: string (nullable = true)
 |    |    |-- namespace: string (nullable = true)
 |    |-- tags: array (nullable = true)
 |    |    |-- element: string (containsNull = true)
 |    |-- title: string (nullable = true)
Run Code Online (Sandbox Code Playgroud)

Ram*_*jan 7

从我的理解到您的问题,可以说您有两个数据框

df1 
root
 |-- col1: string (nullable = true)
 |-- col2: integer (nullable = false)
 |-- col3: double (nullable = false)
Run Code Online (Sandbox Code Playgroud)

和

df2
root
 |-- col1: string (nullable = true)
 |-- col2: string (nullable = true)
 |-- col3: double (nullable = false)
Run Code Online (Sandbox Code Playgroud)

您将必须将的所有列df2合并为一struct列,然后选择要连接的struct列和该列。我在这里col1作为加入专栏

import org.apache.spark.sql.functions._
val nestedDF2 = df2.select($"col1", struct(df2.columns.map(col):_*).as("nested_df2"))
Run Code Online (Sandbox Code Playgroud)

然后,最后一步是join(此处默认为inner join)

df1.join(nestedDF2, Seq("col1"))
Run Code Online (Sandbox Code Playgroud)

这应该给你

root
 |-- col1: string (nullable = true)
 |-- col2: integer (nullable = false)
 |-- col3: double (nullable = false)
 |-- nested_df2: struct (nullable = false)
 |    |-- col1: string (nullable = true)
 |    |-- col2: string (nullable = true)
 |    |-- col3: double (nullable = false)
Run Code Online (Sandbox Code Playgroud)

我希望答案是有帮助的