【问题标题】:Dynamically rename multiple columns in PySpark DataFrame在 PySpark DataFrame 中动态重命名多个列
【发布时间】:2017-05-30 00:55:14
【问题描述】:

我在 pyspark 中有一个数据框,它有 15 列。

列名是idnameemp.dnoemp.salstateemp.cityzip.....

现在我想将其中包含'.' 的列名替换为'_'

喜欢'emp.dno''emp_dno'

我想动态地做它

如何在 pyspark 中实现这一点?

【问题讨论】:

    标签: apache-spark dataframe pyspark special-characters


    【解决方案1】:

    你可以使用类似于this great solution from @zero323的东西:

    df.toDF(*(c.replace('.', '_') for c in df.columns))
    

    或者:

    from pyspark.sql.functions import col
    
    replacements = {c:c.replace('.','_') for c in df.columns if '.' in c}
    
    df.select([col(c).alias(replacements.get(c, c)) for c in df.columns])
    

    replacement 字典看起来像:

    {'emp.city': 'emp_city', 'emp.dno': 'emp_dno', 'emp.sal': 'emp_sal'}
    

    更新:

    如果我的数据框在列名中有空格,那么如何替换 '.''_' 的空格

    import re
    
    df.toDF(*(re.sub(r'[\.\s]+', '_', c) for c in df.columns))
    

    【讨论】:

    • @Virureddy,你能按原样发布print(df.columns) 的输出吗?
    • @Virureddy,谢谢! print(replacements) 长什么样子?
    • @Virureddy,请试试这个:df.toDF([c.replace('.', '_') for c in df.columns])
    • @Virureddy,你可以试试这个:df.toDF(c.replace('.', '_') for c in df.columns)df.toDF((c.replace('.', '_') for c in df.columns))df.toDF(*(c.replace('.', '_') for c in df.columns))
    • 非常感谢您的帮助和耐心
    【解决方案2】:

    编写了一个简单快捷的功能供您使用。享受! :)

    def rename_cols(rename_df):
        for column in rename_df.columns:
            new_column = column.replace('.','_')
            rename_df = rename_df.withColumnRenamed(column, new_column)
        return rename_df
    

    【讨论】:

      【解决方案3】:

      最简单的方法如下:

      解释:

      1. 使用 df.columns 获取 pyspark 数据框中的所有列
      2. 创建一个循环遍历步骤 1 中每一列的列表
      3. 列表将输出:col("col.1").alias(c.replace('.',"_")。仅对需要的列执行此操作。替换功能有助于替换任何模式。另外,你可以排除一些被重命名的列
      4. *[list] 将解压 pypsark 中 select 语句的列表

      from pyspark.sql import functions as F (df .select(*[F.col(c).alias(c.replace('.',"_")) for c in df.columns]) .toPandas().head() )

      希望对你有帮助

      【讨论】:

        【解决方案4】:

        MaxU 的回答既好又高效。这篇文章概述了另一种有效的方法,有助于保持代码库的整洁(使用 quinn 库)。

        假设你有以下DataFrame:

        +---+-----+--------+-------+
        | id| name|emp.city|emp.sal|
        +---+-----+--------+-------+
        | 12|  bob|New York|     80|
        | 99|alice| Atlanta|     90|
        +---+-----+--------+-------+
        

        以下是如何将所有列中的点替换为下划线。

        import quinn
        
        def dots_to_underscores(s):
            return s.replace('.', '_')
        actual_df = df.transform(quinn.with_columns_renamed(dots_to_underscores))
        actual_df.show()
        

        这是结果actual_df

        +---+-----+--------+-------+
        | id| name|emp_city|emp_sal|
        +---+-----+--------+-------+
        | 12|  bob|New York|     80|
        | 99|alice| Atlanta|     90|
        +---+-----+--------+-------+
        

        让我们使用explain() 来验证这个函数是否有效地执行:

        actual_df.explain(True)
        

        这是输出的逻辑计划:

        == Parsed Logical Plan ==
        'Project ['id AS id#50, 'name AS name#51, '`emp.city` AS emp_city#52, '`emp.sal` AS emp_sal#53]
        +- LogicalRDD [id#29, name#30, emp.city#31, emp.sal#32], false
        
        == Analyzed Logical Plan ==
        id: string, name: string, emp_city: string, emp_sal: string
        Project [id#29 AS id#50, name#30 AS name#51, emp.city#31 AS emp_city#52, emp.sal#32 AS emp_sal#53]
        +- LogicalRDD [id#29, name#30, emp.city#31, emp.sal#32], false
        
        == Optimized Logical Plan ==
        Project [id#29, name#30, emp.city#31 AS emp_city#52, emp.sal#32 AS emp_sal#53]
        +- LogicalRDD [id#29, name#30, emp.city#31, emp.sal#32], false
        
        == Physical Plan ==
        *(1) Project [id#29, name#30, emp.city#31 AS emp_city#52, emp.sal#32 AS emp_sal#53]
        

        可以看到解析出来的逻辑计划和物理计划几乎是一样的,所以Catalyst优化器不需要做太多的优化工作。它将id AS id#50 转换为id#29,但这并没有太多的工作。

        with_some_columns_renamed 方法生成更高效的解析计划。

        def dots_to_underscores(s):
            return s.replace('.', '_')
        def change_col_name(s):
          return '.' in s
        actual_df = df.transform(quinn.with_some_columns_renamed(dots_to_underscores, change_col_name))
        actual_df.explain(True)
        

        这个解析后的计划只用点来别名列。

        == Parsed Logical Plan ==
        'Project [unresolvedalias('id, None), unresolvedalias('name, None), '`emp.city` AS emp_city#42, '`emp.sal` AS emp_sal#43]
        +- LogicalRDD [id#34, name#35, emp.city#36, emp.sal#37], false
        
        == Analyzed Logical Plan ==
        id: string, name: string, emp_city: string, emp_sal: string
        Project [id#34, name#35, emp.city#36 AS emp_city#42, emp.sal#37 AS emp_sal#43]
        +- LogicalRDD [id#34, name#35, emp.city#36, emp.sal#37], false
        
        == Optimized Logical Plan ==
        Project [id#34, name#35, emp.city#36 AS emp_city#42, emp.sal#37 AS emp_sal#43]
        +- LogicalRDD [id#34, name#35, emp.city#36, emp.sal#37], false
        
        == Physical Plan ==
        *(1) Project [id#34, name#35, emp.city#36 AS emp_city#42, emp.sal#37 AS emp_sal#43]
        

        更多信息为什么循环遍历 DataFrame 并多次调用 withColumnRenamed 会创建过于复杂的解析计划,应该避免。

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 2015-07-11
          • 2017-01-29
          • 2016-09-23
          • 1970-01-01
          • 2022-01-04
          • 1970-01-01
          • 1970-01-01
          • 2016-12-12
          相关资源
          最近更新 更多