【问题标题】:Adding a new column in Data Frame derived from other columns (Spark)在从其他列(Spark)派生的数据框中添加新列
【发布时间】:2015-09-28 18:28:57
【问题描述】:

我正在使用 Spark 1.3.0 和 Python。我有一个数据框,我希望添加一个从其他列派生的附加列。像这样,

>>old_df.columns
[col_1, col_2, ..., col_m]

>>new_df.columns
[col_1, col_2, ..., col_m, col_n]

在哪里

col_n = col_3 - col_4

如何在 PySpark 中执行此操作?

【问题讨论】:

    标签: python apache-spark apache-spark-sql pyspark


    【解决方案1】:

    实现这一目标的一种方法是使用withColumn 方法:

    old_df = sqlContext.createDataFrame(sc.parallelize(
        [(0, 1), (1, 3), (2, 5)]), ('col_1', 'col_2'))
    
    new_df = old_df.withColumn('col_n', old_df.col_1 - old_df.col_2)
    

    您也可以在已注册的表上使用 SQL:

    old_df.registerTempTable('old_df')
    new_df = sqlContext.sql('SELECT *, col_1 - col_2 AS col_n FROM old_df')
    

    【讨论】:

    • 嘿@zero323,如果我想创建一个列,即 Col_1 是字符串,col_2 是字符串,我希望 column_n 作为 col_1 和 Col_2 的连接,该怎么办?即 Col_1 为零,column_2 为 323。Column_n 应该是 zero323 吗?
    • 谢谢@zero323。虽然我有这个问题: df.select(concat(col("k"), lit(" "), col("v"))) 如何在这里创建第三列?
    • df.withColumn ( 'Datetime', df.column1 + '' + df.column2 ) 不起作用
    • 它不会。你已经有了所有的谜题,你只是把它们放在了正确的地方。
    • 谢谢。使用 df = df.select("*", concat(col("col1"), lit(" "), col("col2")).alias('coln')) 让它工作
    【解决方案2】:

    另外,我们可以使用udf

    from pyspark.sql.functions import udf,col
    from pyspark.sql.types import IntegerType
    from pyspark import SparkContext
    from pyspark.sql import SQLContext
    
    sc = SparkContext()
    sqlContext = SQLContext(sc)
    old_df = sqlContext.createDataFrame(sc.parallelize(
        [(0, 1), (1, 3), (2, 5)]), ('col_1', 'col_2'))
    function = udf(lambda col1, col2 : col1-col2, IntegerType())
    new_df = old_df.withColumn('col_n',function(col('col_1'), col('col_2')))
    new_df.show()
    

    【讨论】:

      【解决方案3】:

      您可以通过以下方式添加新列:

      from pyspark.sql import SparkSession
      
      spark = SparkSession.builder.getOrCreate()
      
      df = spark.createDataFrame([[1, 2], [3, 4]], ['col1', 'col2'])
      df.show()
      
      +----+----+
      |col1|col2|
      +----+----+
      |   1|   2|
      |   3|   4|
      +----+----+
      

      -- 使用withColumn的方法:

      import pyspark.sql.functions as F
      
      df.withColumn('col3', F.col('col2') - F.col('col1')) # col function
      
      df.withColumn('col3', df['col2'] - df['col1']) # bracket notation
      
      df.withColumn('col3', df.col2 - df.col1) # dot notation
      

      -- 使用select的方法:

      df.select('*', (F.col('col2') - F.col('col1')).alias('col3'))
      

      表达式'*' 返回所有列。

      --使用selectExpr的方法:

      df.selectExpr('*', 'col2 - col1 as col3')
      

      -- 使用 SQL:

      df.createOrReplaceTempView('df_view')
      
      spark.sql('select *, col2 - col1 as col3 from df_view')
      

      结果:

      +----+----+----+
      |col1|col2|col3|
      +----+----+----+
      |   1|   2|   1|
      |   3|   4|   1|
      +----+----+----+
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2018-05-08
        • 1970-01-01
        • 2019-04-04
        • 1970-01-01
        • 1970-01-01
        • 2023-03-16
        • 1970-01-01
        相关资源
        最近更新 更多