如何基于多列并按最新Timestamp条件实现PySpark/SQL表关联
关联账户表并处理多环节转账的标签映射(PySpark 3.3.0)
需求说明
- 大表(转账流水)包含字段:
id(唯一发起者ID)、Timestamp(时间戳)、SenderAccount(发起账户)、ReceiverAccount(接收账户) - 小表(账户标签)包含字段:
Account(账户编码)、Label(账户标签),账户编码可匹配大表的发起或接收账户(注意部分账户编码带*,需统一处理匹配) - 核心需求:
- 关联两表,新增
SenderAccLabel(发起账户标签)、ReceiverAccLabel(接收账户标签) - 特殊规则:同一
id对应多笔转账(如A→B→C)时,该id下所有记录的ReceiverAccLabel需取该id下最新Timestamp对应的最终接收账户的标签
- 关联两表,新增
当前表结构
转账流水表(大表)
+---+-------------------+----------------------+----------------------+ |id |Timestamp |SenderAccount |ReceiverAccount | +---+-------------------+----------------------+----------------------+ |1 |2022-01-01 00:00:01|A394840 |*A290375* | |2 |2022-01-01 00:00:01|A943200 |A176108 | |3 |2022-01-01 00:00:02|A139480 |A091392 | |1 |2022-01-01 00:00:02|*A290375* |*A293491* | |4 |2022-01-01 00:00:03|A109347 |A123948 | +---+-------------------+----------------------+----------------------+
账户标签表(小表)
+------------------+-----+ |Account |Label| +------------------+-----+ |A394840 |A3 | |A176108 |B7 | |A139480 |B9 | |*A293491* |*A2* | |A290375 |B6 | +------------------+-----+
期望结果
+---+-------------------+-------------+---------------+--------------+----------------+ |id |Timestamp |SenderAccount|ReceiverAccount|SenderAccLabel|ReceiverAccLabel| +---+-------------------+-------------+---------------+--------------+----------------+ |1 |2022-01-01 00:00:01|A394840 |A290375 |A3 |*A2* | |2 |2022-01-01 00:00:01|A943200 |A176108 |A9 |B7 | |3 |2022-01-01 00:00:02|A139480 |A091392 |B9 |A2 | |1 |2022-01-01 00:00:02|A234870 |A293491 |B7 |*A2* | |4 |2022-01-01 00:00:03|A109347 |A123948 |B6 |A8 | +---+-------------------+-------------+---------------+--------------+----------------+
解决方案
方式一:PySpark API实现
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("AccountLabelMapping").getOrCreate() # ---------------------- # 1. 创建示例数据(实际使用时替换为读取你的数据源) # ---------------------- transactions_data = [ (1, "2022-01-01 00:00:01", "A394840", "*A290375*"), (2, "2022-01-01 00:00:01", "A943200", "A176108"), (3, "2022-01-01 00:00:02", "A139480", "A091392"), (1, "2022-01-01 00:00:02", "*A290375*", "*A293491*"), (4, "2022-01-01 00:00:03", "A109347", "A123948") ] transactions_df = spark.createDataFrame(transactions_data, ["id", "Timestamp", "SenderAccount", "ReceiverAccount"]) transactions_df = transactions_df.withColumn("Timestamp", F.to_timestamp("Timestamp")) account_labels_data = [ ("A394840", "A3"), ("A176108", "B7"), ("A139480", "B9"), ("*A293491*", "*A2*"), ("A290375", "B6"), ("A091392", "A2"), ("A943200", "A9"), ("A109347", "B6"), ("A123948", "A8"), ("A234870", "B7") ] account_labels_df = spark.createDataFrame(account_labels_data, ["Account", "Label"]) # ---------------------- # 2. 核心处理逻辑 # ---------------------- # 清理账户字段的*,统一格式用于匹配 transactions_clean = transactions_df \ .withColumn("SenderAccClean", F.trim(F.lit("*"), F.col("SenderAccount"))) \ .withColumn("ReceiverAccClean", F.trim(F.lit("*"), F.col("ReceiverAccount"))) account_labels_clean = account_labels_df \ .withColumn("AccountClean", F.trim(F.lit("*"), F.col("Account"))) # 获取每个id的最新时间戳对应的最终接收账户 window = Window.partitionBy("id").orderBy(F.desc("Timestamp")) latest_final_receiver = transactions_clean \ .withColumn("row_num", F.row_number().over(window)) \ .filter(F.col("row_num") == 1) \ .select("id", "ReceiverAccClean") \ .withColumnRenamed("ReceiverAccClean", "FinalReceiver") # 关联最终接收账户到原表 transactions_with_final = transactions_clean.join(latest_final_receiver, on="id", how="left") # 关联获取发起账户标签 transactions_with_sender_label = transactions_with_final.join( account_labels_clean, transactions_with_final.SenderAccClean == account_labels_clean.AccountClean, how="left" ) \ .select("id", "Timestamp", "SenderAccount", "ReceiverAccount", "FinalReceiver", F.col("Label").alias("SenderAccLabel")) # 关联获取最终接收账户的标签(即ReceiverAccLabel) final_result = transactions_with_sender_label.join( account_labels_clean, transactions_with_sender_label.FinalReceiver == account_labels_clean.AccountClean, how="left" ) \ .select( "id", "Timestamp", F.trim(F.lit("*"), F.col("SenderAccount")).alias("SenderAccount"), F.trim(F.lit("*"), F.col("ReceiverAccount")).alias("ReceiverAccount"), "SenderAccLabel", F.col("Label").alias("ReceiverAccLabel") ) # 展示结果 final_result.orderBy("id", "Timestamp").show(truncate=False)
方式二:Spark SQL实现
-- 创建临时视图(实际使用时替换为读取你的数据源) CREATE OR REPLACE TEMP VIEW transactions AS SELECT id, to_timestamp(Timestamp) AS Timestamp, SenderAccount, ReceiverAccount FROM VALUES (1, '2022-01-01 00:00:01', 'A394840', '*A290375*'), (2, '2022-01-01 00:00:01', 'A943200', 'A176108'), (3, '2022-01-01 00:00:02', 'A139480', 'A091392'), (1, '2022-01-01 00:00:02', '*A290375*', '*A293491*'), (4, '2022-01-01 00:00:03', 'A109347', 'A123948') AS t(id, Timestamp, SenderAccount, ReceiverAccount); CREATE OR REPLACE TEMP VIEW account_labels AS SELECT Account, Label FROM VALUES ('A394840', 'A3'), ('A176108', 'B7'), ('A139480', 'B9'), ('*A293491*', '*A2*'), ('A290375', 'B6'), ('A091392', 'A2'), ('A943200', 'A9'), ('A109347', 'B6'), ('A123948', 'A8'), ('A234870', 'B7') AS al(Account, Label); -- 核心查询逻辑 WITH clean_transactions AS ( SELECT id, Timestamp, SenderAccount, ReceiverAccount, trim('*', SenderAccount) AS SenderAccClean, trim('*', ReceiverAccount) AS ReceiverAccClean FROM transactions ), latest_receivers AS ( SELECT id, ReceiverAccClean AS FinalReceiver FROM ( SELECT id, ReceiverAccClean, row_number() OVER (PARTITION BY id ORDER BY Timestamp DESC) AS row_num FROM clean_transactions ) t WHERE row_num = 1 ), transactions_with_final_receiver AS ( SELECT ct.id, ct.Timestamp, ct.SenderAccount, ct.ReceiverAccount, ct.SenderAccClean, lr.FinalReceiver FROM clean_transactions ct JOIN latest_receivers lr ON ct.id = lr.id ), transactions_with_sender_label AS ( SELECT twfr.id, twfr.Timestamp, trim('*', twfr.SenderAccount) AS SenderAccount, trim('*', twfr.ReceiverAccount) AS ReceiverAccount, al.Label AS SenderAccLabel, twfr.FinalReceiver FROM transactions_with_final_receiver twfr LEFT JOIN account_labels al ON trim('*', al.Account) = twfr.SenderAccClean ) SELECT id, Timestamp, SenderAccount, ReceiverAccount, SenderAccLabel, al.Label AS ReceiverAccLabel FROM transactions_with_sender_label twsl LEFT JOIN account_labels al ON trim('*', al.Account) = twsl.FinalReceiver ORDER BY id, Timestamp;
内容的提问来源于stack exchange,提问作者rjpost20
相关产品推荐
相关产品推荐

