【问题标题】:How to update dataframe column value while joinining with other dataframe in pyspark?如何在加入pyspark中的其他数据框时更新数据框列值?
【发布时间】:2022-11-10 03:34:10
【问题描述】:

我有 3 个数据框df1(EMPLOYEE_INFO),df2(DEPARTMENT_INFO),df3(COMPANY_INFO),我想通过加入所有三个数据框来更新 df1 中的列。列的名称是 df1 中的 FLAG_DEPARTMENT。我需要设置 FLAG_DEPARTMENT='POLITICS' 。在 sql 查询中将如下所示。

UPDATE [COMPANY_INFO] INNER JOIN ([DEPARTMENT_INFO] 
INNER JOIN [EMPLOYEE_INFO] ON [DEPARTMENT_INFO].DEPT_ID = [EMPLOYEE_INFO].DEPT_ID)
ON [COMPANY_INFO].[COMPANY_DEPT_ID] = [DEPARTMENT_INFO].[DEP_COMPANYID]
SET EMPLOYEE_INFO.FLAG_DEPARTMENT = "POLITICS";

如果这三个表的列中的值匹配,我需要在我的employee_Info 表中设置我的 FLAG_DEPARTMENT='POLITICS'

我怎样才能在 pyspark 中实现同样的目标。我刚开始学习pyspark没有那么深入的知识?

【问题讨论】:

    标签: apache-spark pyspark azure-databricks azure-synapse


    【解决方案1】:

    您可以使用joins 链和select 在其上。

    假设您有以下 pyspark DataFrames:

    employee_df
    +---------+-------+
    |     Name|dept_id|
    +---------+-------+
    |     John| dept_a|
    |      Liù| dept_b|
    |     Luke| dept_a|
    |  Michail| dept_a|
    |      Noe| dept_e|
    |Shinchaku| dept_c|
    |     Vlad| dept_e|
    +---------+-------+
    
    department_df
    +-------+----------+------------+
    |dept_id|company_id| description|
    +-------+----------+------------+
    | dept_a|  company1|Department A|
    | dept_b|  company2|Department B|
    | dept_c|  company5|Department C|
    | dept_d|  company3|Department D|
    +-------+----------+------------+
    
    company_df
    +----------+-----------+
    |company_id|description|
    +----------+-----------+
    |  company1|  Company 1|
    |  company2|  Company 2|
    |  company3|  Company 3|
    |  company4|  Company 4|
    +----------+-----------+
    

    然后您可以运行以下代码将flag_department 列添加到您的employee_df

    from pyspark.sql import functions as F
    
    employee_df = (
            employee_df.alias('a')
            .join(
                department_df.alias('b'),
                on='dept_id',
                how='left',
            )
            .join(
                company_df.alias('c'),
                on=F.col('b.company_id') == F.col('c.company_id'),
                how='left',
            )
            .select(
                *[F.col(f'a.{c}') for c in employee_df.columns],
                F.when(
                    F.col('b.dept_id').isNotNull() & F.col('c.company_id').isNotNull(),
                    F.lit('POLITICS')
                ).alias('flag_department')
            )
        )
    

    新的employee_df 将是:

    +---------+-------+---------------+
    |     Name|dept_id|flag_department|
    +---------+-------+---------------+
    |     John| dept_a|       POLITICS|
    |      Liù| dept_b|       POLITICS|
    |     Luke| dept_a|       POLITICS|
    |  Michail| dept_a|       POLITICS|
    |      Noe| dept_e|           null|
    |Shinchaku| dept_c|           null|
    |     Vlad| dept_e|           null|
    +---------+-------+---------------+
    

    【讨论】:

    • 嗨@PieCot 如果在employee_df 列名是DEPT_ID 并且在Department_df 列名是DEPARTMENT_IDS 会发生什么
    • 在这种情况下,我们不能直接加入
    • 您可以在第一次加入时更改子句onon=F.col('a.DEPT_ID') == F.col('b.DEPARTMENT_IDS')
    • 嗨,兄弟,如果我删除第二个联接并仅将 select 与第一个联接一起使用,那么它也适用于更新 empdf
    • 完美的!根据您想要更新的条件,两者都可以工作:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多