Pyspark比较数据框中的记录并查找每列的增量

如何解决Pyspark比较数据框中的记录并查找每列的增量

我有两个数据帧,fulldata和deltadata,需要进行比较,并需要根据以下规则进行输出

  1. 如果两个列的值都相同,则不应在输出中显示-例如emp1
  2. 当增量数据发生变化时,输出中应包含已更改的列以及在MANDATORY_FIELDS_LOOKUP中定义的强制列–例如emp2,其中已更改的列“ field2”和强制列employee_number,record_type,姓,first_forename和date_of_birth已添加到输出中(字段1和字段3应该为空,因为没有变化)
  3. 如果更改了必填列,则应包括该记录-例如emp3,其中更改的列是姓和first_forename
  4. 应该包括增量中的新记录–例如emp5
  5. 任何更改为空字符串的列都应硬编码为“(空白)” –例如emp6

我有一个工作代码可以提供所需的输出,但是性能却很慢– 250条记录的500条记录大约需要10分钟。我正在使用以下配置在连接到Glue Dev Endpoint的Juypter Notebook中运行此代码。 工人人数15 工种G.2X 数据处理单元(DPU)31 能否以更有效的方式实施?任何帮助表示赞赏。

from pyspark.sql.types import StructType,StringType,StructField
from pyspark.context import SparkContext
from pyspark.sql import SparkSession,DataFrame
from pyspark import SQLContext
from pyspark.sql.functions import *
from pyspark.sql.window import Window
from operator import add
from functools import reduce

spark = SparkContext.getOrCreate()
sql_context = SQLContext(spark)
spark_session = SparkSession.builder.enableHiveSupport().getOrCreate()

MANDATORY_FIELDS_LOOKUP = {
    "10PERDET": [
        "surname","first_forename","date_of_birth",]
}

def delta_compare(full_df: DataFrame,delta_df: DataFrame,record_header:str,join_key: list) -> DataFrame:
    all_fields = full_df.columns
    mandatory_fields = MANDATORY_FIELDS_LOOKUP.get(record_header)
    delta_df = delta_df.withColumn("isDelta",lit(1))
    full_df = full_df.withColumn("isDelta",lit(0))
    union_df = full_df.union(delta_df)

    window_spec = Window.partitionBy([col(k) for k in join_key]).orderBy(col("isDelta"))

    df = (
        union_df.withColumn(
            "isChanged",reduce(
                add,[
                    when(col(fld) == lag(col(fld)).over(window_spec),0).otherwise(1)
                    for fld in all_fields
                    if fld not in ["record_type",*join_key,"isDelta"]
                ],),)
        .select(
            *join_key,col("record_type"),*mandatory_fields,*[
                when(col(fld) == lag(col(fld)).over(window_spec),"")
                .when(
                    (col(fld) != lag(col(fld)).over(window_spec))
                    & (col(fld) == ""),"(blank)"
                )
                .otherwise(col(fld))
                .alias(fld)
                for fld in all_fields
                if fld not in ["record_type","isDelta"] + mandatory_fields
            ],)
        .filter(col("isDelta") == 1)
        .filter(col("isChanged") != 0)
    )
    return df.select(*[fld for fld in all_fields]).orderBy(*join_key)

schema = (StructType([StructField("employee_number",StringType()),StructField("record_type",StructField("surname",StructField("first_forename",StructField("date_of_birth",StructField("field1",StructField("field2",StructField("field3",StringType())])
         )
fulldata = [['emp1',"10PERDET",'Neil','Par','16011980','10','20','30'],['emp2','Tom','Hanks','11091982','15','25','35'],['emp3','jag','ram','26121981','17','27','37'],['emp4','right','sam','26121990','oldrow','99','88'],['emp6','coke','john','01021985','29','39','49'],]
full_df = sql_context.createDataFrame(fulldata,schema=schema)

deltadata = [['emp1','new','jag_new','ram_new',['emp5','newjohn','gan','22022020','newrow','01','02'],'',]
delta_df = sql_context.createDataFrame(deltadata,schema=schema)
final_df = delta_compare(full_df,delta_df,['employee_number'])
final_df.show()

输出:

+---------------+-----------+-------+--------------+-------------+------+-------+------+
|employee_number|record_type|surname|first_forename|date_of_birth|field1| field2|field3|
+---------------+-----------+-------+--------------+-------------+------+-------+------+
|           emp2|   10PERDET|    Tom|         Hanks|     11091982|      |    new|      |
|           emp3|   10PERDET|jag_new|       ram_new|     26121981|   new|       |      |
|           emp5|   10PERDET|newjohn|           gan|     22022020|newrow|     01|    02|
|           emp6|   10PERDET|   coke|          john|     01021985|      |(blank)|      |
+---------------+-----------+-------+--------------+-------------+------+-------+------+

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


依赖报错 idea导入项目后依赖报错,解决方案:https://blog.csdn.net/weixin_42420249/article/details/81191861 依赖版本报错:更换其他版本 无法下载依赖可参考:https://blog.csdn.net/weixin_42628809/a
错误1:代码生成器依赖和mybatis依赖冲突 启动项目时报错如下 2021-12-03 13:33:33.927 ERROR 7228 [ main] o.s.b.d.LoggingFailureAnalysisReporter : *************************** APPL
错误1:gradle项目控制台输出为乱码 # 解决方案:https://blog.csdn.net/weixin_43501566/article/details/112482302 # 在gradle-wrapper.properties 添加以下内容 org.gradle.jvmargs=-Df
错误还原:在查询的过程中,传入的workType为0时,该条件不起作用 <select id="xxx"> SELECT di.id, di.name, di.work_type, di.updated... <where> <if test=&qu
报错如下,gcc版本太低 ^ server.c:5346:31: 错误:‘struct redisServer’没有名为‘server_cpulist’的成员 redisSetCpuAffinity(server.server_cpulist); ^ server.c: 在函数‘hasActiveC
解决方案1 1、改项目中.idea/workspace.xml配置文件,增加dynamic.classpath参数 2、搜索PropertiesComponent,添加如下 <property name="dynamic.classpath" value="tru
删除根组件app.vue中的默认代码后报错:Module Error (from ./node_modules/eslint-loader/index.js): 解决方案:关闭ESlint代码检测,在项目根目录创建vue.config.js,在文件中添加 module.exports = { lin
查看spark默认的python版本 [root@master day27]# pyspark /home/software/spark-2.3.4-bin-hadoop2.7/conf/spark-env.sh: line 2: /usr/local/hadoop/bin/hadoop: No s
使用本地python环境可以成功执行 import pandas as pd import matplotlib.pyplot as plt # 设置字体 plt.rcParams['font.sans-serif'] = ['SimHei'] # 能正确显示负号 p
错误1:Request method ‘DELETE‘ not supported 错误还原:controller层有一个接口,访问该接口时报错:Request method ‘DELETE‘ not supported 错误原因:没有接收到前端传入的参数,修改为如下 参考 错误2:cannot r
错误1:启动docker镜像时报错:Error response from daemon: driver failed programming external connectivity on endpoint quirky_allen 解决方法:重启docker -> systemctl r
错误1:private field ‘xxx‘ is never assigned 按Altʾnter快捷键,选择第2项 参考:https://blog.csdn.net/shi_hong_fei_hei/article/details/88814070 错误2:启动时报错,不能找到主启动类 #
报错如下,通过源不能下载,最后警告pip需升级版本 Requirement already satisfied: pip in c:\users\ychen\appdata\local\programs\python\python310\lib\site-packages (22.0.4) Coll
错误1:maven打包报错 错误还原:使用maven打包项目时报错如下 [ERROR] Failed to execute goal org.apache.maven.plugins:maven-resources-plugin:3.2.0:resources (default-resources)
错误1:服务调用时报错 服务消费者模块assess通过openFeign调用服务提供者模块hires 如下为服务提供者模块hires的控制层接口 @RestController @RequestMapping("/hires") public class FeignControl
错误1:运行项目后报如下错误 解决方案 报错2:Failed to execute goal org.apache.maven.plugins:maven-compiler-plugin:3.8.1:compile (default-compile) on project sb 解决方案:在pom.
参考 错误原因 过滤器或拦截器在生效时,redisTemplate还没有注入 解决方案:在注入容器时就生效 @Component //项目运行时就注入Spring容器 public class RedisBean { @Resource private RedisTemplate<String
使用vite构建项目报错 C:\Users\ychen\work>npm init @vitejs/app @vitejs/create-app is deprecated, use npm init vite instead C:\Users\ychen\AppData\Local\npm-