【问题标题】:How can I update Pyspark DataFrame column values under two column conditions using Bitwise or bit and function?如何使用按位或位和函数在两列条件下更新 Pyspark DataFrame 列值?
【发布时间】:2022-07-02 01:09:11
【问题描述】:

我需要在pyspark数据框中更新一列(Flag,包含许多标志,每个标志是2^nint数字,加起来)在两个条件下,即column(Age)值> = 65 和列 Flag 不包含由按位或位和函数检查的新标志值:(Flag & newFlag) == 0

我已经使用示例数据框和 python 脚本演示了我的工作(请参见下文),但遇到了错误消息。 错误信息是:AnalysisException: cannot resolve '(Flag AND 2)' due to data type mismatch: '(Flag AND 2)' requires boolean type, not int;

from pyspark.sql.types import StructType,StructField, StringType, IntegerType`
from pyspark.sql.functions import *

# create a data frame with two columns: Age and Flag and three rows
data = [
(61,0),
(65,1),
(66,10)  #previous inserted Flag 2 and 8, add up to 10, Flag is 2^n
]
schema = StructType([ \
StructField("Age",IntegerType(), True), \
StructField("Flag",IntegerType(), True) \
])

df = spark.createDataFrame(data=data,schema=schema)
#df.printSchema()
df.show(truncate=False)

N_FLAG_AGE65=2
new_column = when(
   (col("Age") >= 65) & ((col("Flag") & lit(N_FLAG_AGE65) == 0)), 
   col("Flag")+N_FLAG_AGE65     
).otherwise(col("Flag"))
df = df.withColumn("Flag", new_column)
df.show(truncate=False)

【问题讨论】:

  • 请添加您的示例输入和预期输出数据集。它将让论坛以更好的方式了解您的用例。

标签: dataframe pyspark bitwise-operators


【解决方案1】:

输入源df构造完成后,df.show(truncate=False)的第一显示行应该是

+---+----+ |年龄|国旗| +---+----+ |61 |0 | |65 |1 | |66 |10 | +---+----+

我的更新算法是检查两列(Age 和 Flag),如果 age >=65 并且 Flag 位函数不包含 N_FLAG_AGE65,我们将 Flag 字段更新为 Flag = Flag+N_FLAG_AGE65。因此,预期的结果应该是 +---+----+ |年龄|国旗| +---+----+ |61 |0 | |65 |3 | |66 |10 | +---+----+

我认为“new_column”条件表达式的原始语法不适用于 df = df.withColumn("Flag", new_column)

我做了语法更改,它现在适用于以下代码,方法是添加一个名为 column(Flag65_exp) 的新常量 lit(N_FLAG_AGE65) 并使用 expr("case when Age>=65 and Flag & lit(N_FLAG_AGE65)=0 then Flag+lit(N_FLAG_AGE65) Else Flag End") in df.withColumn("Flag",expr("..."))

%python
from pyspark.sql.types import StructType,StructField, 
StringType, IntegerType
from pyspark.sql.functions import *

# create a data frame with two columns: Age and Flag and three 
rows
data = [
(61,0),
(65,1),
(66,10)  #previous inserted Flag 2 and 8, add up to 10, Flag is 
2^n
]
schema = StructType([ \
StructField("Age",IntegerType(), True), \
StructField("Flag",IntegerType(), True) \
])

df = spark.createDataFrame(data=data,schema=schema)
#df.printSchema()
df.show(truncate=False)

N_FLAG_AGE65=2
df=df.withColumn('Flag65_exp', lit(N_FLAG_AGE65))
df = df.withColumn("Flag", expr("case when Age>=65 and Flag & 
lit(N_FLAG_AGE65)=0 then Flag+lit(N_FLAG_AGE65) Else Flag End"))
df.show(truncate=False)

#source df +---+----+ |年龄|国旗| +---+----+ |61 |0 | |65 |1 | |66 |10 | +---+----+

#更新的df +---+----+----------+ |年龄|标志|Flag65_exp| +---+----+----------+ |61 |0 |2 | |65 |3 |2 | |66 |10 |2 | +---+----+----------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-11-23
    • 2017-03-16
    • 1970-01-01
    • 2020-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多