【问题标题】:Text file comparison using Spark data frames使用 Spark 数据帧的文本文件比较
【发布时间】:2018-10-08 18:04:45
【问题描述】:

我想使用 Spark 数据帧来实现以下要求来比较 2 个文本/csv

  • 列表项

文件。理想情况下,File1.txt 应该与 File2.txt 进行比较,结果应该在其他带有标志的 txt 文件中(SAME/UPDATE/INSERT/DELETE)。

UPDATE - 与 file1 相比,如果 file2 中的任何记录值更新 INSERT - 如果 file2 中存在新记录 DELETE - 仅当记录存在于 file1 中(不在 file2 中) SAME - 如果两个文件中存在相同的记录

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

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

Outputfile.txt
NO  DEPT NAME    SAL   FLAG
1   IT  RAM     1000    S
2   IT  SRI     900     U
4   MT  SUMP    1200    I
5   HW  MAHI    700     S
3   HR  GOPI    1500    D

到目前为止,我做了以下编码。但无法进一步进行。请帮忙。

from pyspark.shell import spark
sc = spark.sparkContext
df1 = spark.read.option("header","true").option("delimiter", ",").csv("C:\\inputs\\file1.csv")
df2 = spark.read.option("header","true").option("delimiter", ",").csv("C:\\inputs\\file2.csv")

df1.createOrReplaceTempView("table1")
df2.createOrReplaceTempView("table2")

sqlDF1 = spark.sql( "select * from table1" )
sqlDF2 = spark.sql( "select * from table2" )

leftJoinDF = sqlDF1.join(sqlDF2, 'id', how='left')
rightJoinDF = sqlDF1.join(sqlDF2, 'id', how='right')
innerJoinDF = sqlDF1.join(sqlDF2, 'id')

如果我们在执行leftJoin,rightJoin,innerJoin之后合并数据,有什么办法。有了这个,我是否可以获得所需的输出或任何其他方式。

谢谢,

【问题讨论】:

    标签: pyspark apache-spark-sql


    【解决方案1】:

    您可以在下面找到我的解决方案。我为 SAME/UPDATE/INSERT/DELETE 案例创建 4 个数据框,然后合并它们

    >>> from functools import reduce
    >>> from pyspark.sql import DataFrame
    >>> import pyspark.sql.functions as F
    
    >>> df1 = sc.parallelize([
    ...     (1,'IT','RAM',1000),    
    ...     (2,'IT','SRI',600),
    ...     (3,'HR','GOPI',1500),    
    ...     (5,'HW','MAHI',700)
    ...     ]).toDF(['NO','DEPT','NAME','SAL'])
    >>> df1.show()
    +---+----+----+----+
    | NO|DEPT|NAME| SAL|
    +---+----+----+----+
    |  1|  IT| RAM|1000|
    |  2|  IT| SRI| 600|
    |  3|  HR|GOPI|1500|
    |  5|  HW|MAHI| 700|
    +---+----+----+----+
    
    >>> df2 = sc.parallelize([
    ...     (1,'IT','RAM',1000),    
    ...     (2,'IT','SRI',900),
    ...     (4,'MT','SUMP',1200),    
    ...     (5,'HW','MAHI',700)
    ...     ]).toDF(['NO','DEPT','NAME','SAL'])
    >>> df2.show()
    +---+----+----+----+
    | NO|DEPT|NAME| SAL|
    +---+----+----+----+
    |  1|  IT| RAM|1000|
    |  2|  IT| SRI| 900|
    |  4|  MT|SUMP|1200|
    |  5|  HW|MAHI| 700|
    +---+----+----+----+
    
    #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('D').alias('FLAG'))
    #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('I').alias('FLAG'))
    #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('S').alias('FLAG'))
    #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('U').alias('FLAG'))
    
    
    >>> dfs = [df_s,df_u,df_u,df_i]
    >>> df = reduce(DataFrame.unionAll, dfs)
    >>> 
    >>> df.show()
    +---+----+----+----+----+                                                       
    | NO|DEPT|NAME| SAL|FLAG|
    +---+----+----+----+----+
    |  5|  HW|MAHI| 700|   S|
    |  1|  IT| RAM|1000|   S|
    |  2|  IT| SRI| 900|   U|
    |  2|  IT| SRI| 900|   U|
    |  4|  MT|SUMP|1200|   I|
    +---+----+----+----+----+
    

    【讨论】:

    • 感谢分享,但在上面的最终结果中我想确认2个澄清:1.记录(NO:2),重复两次。只有数据正确。
    • 对不起,我用了 df_u 两次。你应该用 df_d 改变第二个 df_u。它应该是 dfs = [df_s,df_u,df_d,df_i]。我也用过 spark 2.3.0
    • 试试这个 -- df = df_d.union(df_i).union(df_s).union(df_u)
    • 在 df_s 和 df_u 中的 concat 期间,如果 (no, dept, name, sal) 中有任何空值,则该记录不会出现在最终输出中。有什么方法可以过滤空值并将它们包含在合适的标志下。
    • 我尝试添加如下的空条件: (df2.NO!='null', df2.DEPT!='null', df2.NAME!='null' ,df2.SAL! ='null') ,但问题仍然存在。
    【解决方案2】:

    您可以在先连接所有列后使用'outer' join。然后为标志创建一个udf

    import pyspark.sql.functions as F
    
    df = sql.createDataFrame([
         (1,'IT','RAM',1000),
         (2,'IT','SRI',600),
         (3,'HR','GOPI',1500),
         (5,'HW','MAHI',700)],
         ['NO'  ,'DEPT', 'NAME',   'SAL' ])
    
    df1 = sql.createDataFrame([
         (1,'IT','RAM',1000),
         (2,'IT','SRI',900),
         (4,'MT','SUMP',1200 ),
         (5,'HW','MAHI',700)],
         ['NO'  ,'DEPT', 'NAME',   'SAL' ])
    
    def flags(x,y):
        if not x:
            return y+'-I'
        if not y:
            return x+'-D'
        if x == y:
            return x+'-S'
        return y+'-U'
    
    _cols = df.columns
    flag_udf = F.udf(lambda x,y: flags(x,y),StringType())   
    
    
    df = df.select(['NO']+ [F.concat_ws('-', *[F.col(_c) for _c in df.columns]).alias('f1')])\
            .join(df1.select(['NO']+ [F.concat_ws('-', *[F.col(_c1) for _c1 in df1.columns]).alias('f2')]), 'NO', 'outer')\
            .select(flag_udf('f1','f2').alias('combined'))
    df.show()
    

    结果会是,

    +----------------+                                                              
    |        combined|
    +----------------+
    | 5-HW-MAHI-700-S|
    | 1-IT-RAM-1000-S|
    |3-HR-GOPI-1500-D|
    |  2-IT-SRI-900-U|
    |4-MT-SUMP-1200-I|
    +----------------+
    

    最后,拆分combined 列。

    split_col = F.split(df['combined'], '-')
    df = df.select([split_col.getItem(i).alias(s) for i,s in enumerate(_cols+['FLAG'])])
    
    df.show()
    

    你得到想要的输出,

    +---+----+----+----+----+                                                       
    | NO|DEPT|NAME| SAL|FLAG|
    +---+----+----+----+----+
    |  5|  HW|MAHI| 700|   S|
    |  1|  IT| RAM|1000|   S|
    |  3|  HR|GOPI|1500|   D|
    |  2|  IT| SRI| 900|   U|
    |  4|  MT|SUMP|1200|   I|
    +---+----+----+----+----+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-01-15
      • 2020-06-07
      • 1970-01-01
      • 2019-06-02
      • 1970-01-01
      • 2021-07-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多