【发布时间】: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