【问题标题】:Databricks - How to join a table with IDs contained in a column of type struct<array<string>>Databricks - 如何使用 struct<array<string>> 类型的列中包含的 ID 连接表
【发布时间】:2021-12-06 11:05:35
【问题描述】:

我目前有 JSON 文件,我可以通过该文件将其数据转储到临时视图中。遵循 Python (PySpark) 逻辑:

 departMentData = spark \
                .read \
                .option("multiLine", True) \
                .option("mode", "PERMISSIVE") \
                .json("C:\\Test\data.json") \
                .createOrReplaceTempView("vw_TestView")

此临时视图以数组的形式包含部门数据和该部门内的员工列表。一名员工可以隶属于多个部门。

以下是该视图的数据类型:

  • 部门ID:字符串
  • 部门名称:字符串
  • EmployeeIDs:数组

vw_TestView的表格数据如下

DeptID DeptName EmployeeIDs
D01 dev ["U1234", "U6789"]
D02 qa ["U1234", "U2345"]

另一个表Employees包含所有这些员工的详细信息,如下所示:

EmpID EmpName
U1234 jon
U6789 smith
U2345 natasha

我需要一个新表的最终输出如下:

DeptID DeptName EmployeeIDs EmployeeNames
D01 dev ["U1234", "U6789"] ["jon", "smith"]
D02 qa ["U1234", "U2345"] ["jon", "natasha"]

如何在 Databricks SQL 中或通过 PySPark 执行此类连接?

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql databricks delta-lake


    【解决方案1】:

    您可以尝试以下方法,它使用explode 将员工 ID 列表拆分为不同的行,然后再加入它们并使用collect_list 将条目聚合到一个列表中。

    使用 spark sql:

    注意。 确保Employees 可用作表格/视图,例如EmployeeData.createOrReplaceTempView("Employees")

    WITH dept_employees AS (
        SELECT
            DeptId,
            DeptName,
            explode(EmployeeIDs)
        FROM
            vw_TestView
    )
    SELECT
        d.DeptId,
        d.DeptName,
        collect_list(e.EmpID) as EmployeeIDs,
        collect_list(e.EmpName) as EmployeeNames
    FROM
        dept_employees d
    INNER JOIN
        Employees e ON d.col=e.EmpID
    GROUP BY
        d.Deptid,
        d.DeptName
    

    或使用 pyspark api:

    from pyspark.sql import functions as F
    
    output_df = (
        departMentData.select(
            F.col("DeptId"),
            F.col("DeptName"),
            F.explode("EmployeeIDs")
        )
        .alias("d")
        .join(
            EmployeeData.alias("e"),
            F.col("d.col")==F.col("e.EmpID"),
            "inner"
        )
        .groupBy("d.DeptId","d.DeptName")
        .agg(
            F.collect_list("e.EmpID").alias("EmployeeIDs"),
            F.collect_list("e.EmpName").alias("EmployeeNames")
        )
    )
    

    让我知道这是否适合你。

    【讨论】:

    • 这两种方法都有效。
    猜你喜欢
    • 2015-09-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-22
    • 2018-08-16
    • 2022-12-17
    • 1970-01-01
    相关资源
    最近更新 更多