PySpark无需staging表实现相同customerID的列值累加更新
全量增量DataFrame金额累加合并方案
核心实现逻辑无需创建任何临时表,仅通过一次外连接+空值补0求和即可覆盖所有需求场景:
- 存量未出现在增量中的客户:增量金额补0,求和结果为原存量金额
- 存量增量都存在的客户:两个金额直接累加
- 增量新增客户:存量金额补0,求和结果为增量金额
Pandas 实现代码
import pandas as pd # 全量表 full_df、增量表 incremental_df 为已读入的DataFrame对象 # 1. 按customer_ID做外连接,字段加后缀区分来源 merge_temp = full_df.merge( incremental_df, on="customer_ID", how="outer", suffixes=("_full", "_inc") ) # 2. 空值填充为0后求和得到最终amount merge_temp["amount"] = merge_temp["amount_full"].fillna(0) + merge_temp["amount_inc"].fillna(0) # 3. 提取目标字段得到结果 result = merge_temp[["customer_ID", "amount"]]
PySpark 实现代码
from pyspark.sql import functions as F # 先给增量表的amount字段重命名,避免连接后字段重名冲突 inc_rename = incremental_df.withColumnRenamed("amount", "inc_amount") # 1. 按customer_ID做外连接 merge_temp = full_df.join(inc_rename, on="customer_ID", how="outer") # 2. 空值补0后求和得到最终amount result = merge_temp.select( "customer_ID", (F.coalesce(F.col("amount"), F.lit(0)) + F.coalesce(F.col("inc_amount"), F.lit(0))).alias("amount") )
内容的提问来源于stack exchange,提问作者bigDataArtist
相关产品推荐
相关产品推荐

