【问题标题】:How to check NULL values while comparing 2 text files using spark data frames如何在使用火花数据帧比较 2 个文本文件时检查 NULL 值
【发布时间】:2018-10-10 10:12:36
【问题描述】:

以下代码未能捕获“空”值记录。从 df1 下方,NO 列。 5 有一个空值(名称字段)。

按照我下面的 OutputDF 要求,第 5 名的记录应该如前所述。但是在下面的代码执行之后,这条记录不会进入最终输出。具有“空”值的记录不会进入输出。除此之外,一切都很好。

df1

NO  DEPT NAME   SAL 
1   IT  RAM     1000    
2   IT  SRI     600 
3   HR  GOPI    1500    
5   HW          700

df2

NO  DEPT NAME   SAL 
1   IT   RAM    1000    
2   IT   SRI    900 
4   MT   SUMP   1200    
5   HW   MAHI   700

输出DF

NO  DEPT NAME    SAL   FLAG
1   IT  RAM     1000   SAME
2   IT  SRI     900    UPDATE
4   MT  SUMP    1200   INSERT
3   HR  GOPI    1500   DELETE
5   HW  MAHI    700    UPDATE

from pyspark.shell import spark
from pyspark.sql import DataFrame
import pyspark.sql.functions as F
sc = spark.sparkContext

filedf1 = spark.read.option("header","true").option("delimiter", ",").csv("C:\\files\\file1.csv")
filedf2 = spark.read.option("header","true").option("delimiter", ",").csv("C:\\files\\file2.csv")
filedf1.createOrReplaceTempView("table1")
filedf2.createOrReplaceTempView("table2")
df1 = spark.sql( "select * from table1" )
df2 = spark.sql( "select * from table2" )

#DELETE
df_d = df1.join(df2, df1.NO == df2.NO, 'left').filter(F.isnull(df2.NO)).select(df1.NO,df1.DEPT,df1.NAME,df1.SAL, F.lit('DELETE').alias('FLAG'))
print("df_d left:",df_d.show())
#INSERT
df_i = df1.join(df2, df1.NO == df2.NO, 'right').filter(F.isnull(df1.NO)).select(df2.NO,df2.DEPT,df2.NAME,df2.SAL, F.lit('INSERT').alias('FLAG'))
print("df_i right:",df_i.show())
#SAME
df_s = df1.join(df2, df1.NO == df2.NO, 'inner').filter(F.concat(df2.NO,df2.DEPT,df2.NAME,df2.SAL) == F.concat(df1.NO,df1.DEPT,df1.NAME,df1.SAL)).select(df1.NO,df1.DEPT,df1.NAME,df1.SAL, F.lit('SAME').alias('FLAG'))
print("df_s inner:",df_s.show())
#UPDATE
df_u = df1.join(df2, df1.NO == df2.NO, 'inner').filter(F.concat(df2.NO,df2.DEPT,df2.NAME,df2.SAL) != F.concat(df1.NO,df1.DEPT,df1.NAME,df1.SAL)).select(df2.NO,df2.DEPT,df2.NAME,df2.SAL, F.lit('UPDATE').alias('FLAG'))
print("df_u inner:",df_u.show())

df = df_d.union(df_i).union(df_s).union(df_u)
df.show()

这里我比较 df1 和 df2,如果在 df2 中发现新记录以标志作为 INSERT,如果记录在两个 dfs 中相同,则作为 SAME,如果记录在 df1 中而不在 df2 中,则作为 DELETE 和如果记录在两个 dfs 中都存在但具有不同的值,则将 df2 值作为 UPDATE。

【问题讨论】:

  • 您确定是Null 而不是space(空字符)?
  • 对不起,它不应该只接受空值。在 sql 中,我们将使用类似 nvl(name,'') 的东西,甚至我们也可以在这里使用类似的东西来转义空值。因为如果我们采用大输入文件(csv),我们不知道空值在哪里,这些事情我需要在这里防止。因此,我需要实现这一点。此验证仅待处理,其余代码工作正常。

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


【解决方案1】:

代码有两个问题:

  1. f.concat的结果为null返回null,所以这部分代码过滤掉了第5行:

    .filter(F.concat(df2.NO, df2.NAME, df2.SAL) != F.concat(df1.NO, df1.NAME, df1.SAL))
    
  2. 您只选择了 df2。在上面的示例中很好,但是如果您的 df2 为 null,则生成的数据框将为 null。

您可以尝试将其与下面的 udf 连接:

def concat_cols(row):
    concat_row = ''.join([str(col) for col in row if col is not None])
    return concat_row 

udf_concat_cols = udf(concat_cols, StringType())

函数concat_row可以分为两部分:

  1. "".join([mylist]) 是string function。它加入了一切 带有定义的分隔符的列表,在这种情况下它是一个空字符串。
  2. [str(col) for col in row if col is not None] 是一个列表推导,它的作用是:对于行中的每一列,如果 该列不是 None,然后将 str(col) 附加到列表中。
    List comprehension 只是一种更 Pythonic 的方式:

    mylist = [] 
    for col in row: 
        if col is not None:
            mylist.append(col))
    

您可以将更新代码替换为:

df_u = (df1
.join(df2, df1.NO == df2.NO, 'inner')
.filter(udf_concat_cols(struct(df1.NO, df1.NAME, df1.SAL)) != udf_concat_cols(struct(df2.NO, df2.NAME, df2.SAL)))
.select(coalesce(df1.NO, df2.NO), 
        coalesce(df1.NAME, df2.NAME),
        coalesce(df1.SAL, df2.SAL),
        F.lit('UPDATE').alias('FLAG')))

您应该为#SAME 标志做类似的事情并换行以提高可读性。


更新:

如果 df2 始终具有正确(更新)的结果,则无需合并。 这个实例的代码是:

df_u = (df1
.join(df2, df1.NO == df2.NO, 'inner')
.filter(udf_concat_cols(struct(df1.NO, df1.NAME, df1.SAL)) != udf_concat_cols(struct(df2.NO, df2.NAME, df2.SAL)))
.select(df2.NO,
        df2.NAME,
        df2.SAL,
        F.lit('UPDATE').alias('FLAG')))

【讨论】:

  • 非常感谢。问题快结束了。除了以下更改: NO,NAME,DEPT,SAL 4,ABC,,9000 (df1 : DEPT is null) 4,DEF,, (df2 : DEPT, SAL is null) 当前输出:4,DEF,null,9000 (UPDATE ) 例外输出:4,DEF,null,null
  • 当df2中的某个字段为null时,它正在从df1中获取相应的值。其实不应该是这样的。 DF2 是最终的,在 UPDATE 的情况下不应该再去 DF1。
  • 如果是这种情况,那么不要合并,只像以前那样选择 df2 ,那么这会给你正确的结果。检查更新的答案是否正确?
  • 对于较长的函数,需要先定义函数,注册为udf,然后就可以像普通的lambda udf一样使用了。但是,对于像 udf_a = udf(lambda row: row.DoSomething, StringType()) 这样的 lambda udf,您不需要注册它。看看databricks tutorial
猜你喜欢
  • 2017-01-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-09-13
  • 2021-12-14
  • 2020-06-07
  • 1970-01-01
  • 2018-02-15
相关资源
最近更新 更多