You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于多列并按最新Timestamp条件实现PySpark/SQL表关联

关联账户表并处理多环节转账的标签映射(PySpark 3.3.0)

需求说明

  • 大表(转账流水)包含字段:id(唯一发起者ID)、Timestamp(时间戳)、SenderAccount(发起账户)、ReceiverAccount(接收账户)
  • 小表(账户标签)包含字段:Account(账户编码)、Label(账户标签),账户编码可匹配大表的发起或接收账户(注意部分账户编码带*,需统一处理匹配)
  • 核心需求:
    1. 关联两表,新增SenderAccLabel(发起账户标签)、ReceiverAccLabel(接收账户标签)
    2. 特殊规则:同一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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 12:54:23