如何用PySpark基于汇率表转换DataFrame中的货币金额?
PySpark 按年月匹配汇率转换货币值实现方案
核心思路
通过给两个DataFrame提取年月维度的关联键,结合货币对(from_curr、to_curr)进行关联,再用汇率计算转换后的货币值。
具体实现步骤
1. 导入PySpark函数库
from pyspark.sql import functions as F
2. 处理待转换数据df1:提取年月字段
先从Date字段中提取yyyy-MM格式的年月,作为关联的时间维度键:
# 若Date字段是字符串类型,先转换为日期格式:df1 = df1.withColumn("Date", F.to_date("Date", "yyyy-MM-dd")) df1_with_month = df1.withColumn("year_month", F.date_format(F.col("Date"), "yyyy-MM"))
3. 处理汇率表df2:统一字段名+提取年月
注意df2的目标货币字段是To_curr,需和df1的to_curr统一字段名,再提取年月:
# 统一字段名,提取年月维度 df2_clean = df2.withColumnRenamed("To_curr", "to_curr") \ .withColumn("year_month", F.date_format(F.col("Date"), "yyyy-MM"))
4. 清理汇率表:去重保留最新汇率
若汇率表中同一货币对、同一年月存在多条记录,按日期倒序取最新的汇率:
df2_latest_rate = df2_clean.orderBy("Date", ascending=False) \ .groupBy("from_curr", "to_curr", "year_month") \ .agg(F.first("rate_exchange").alias("latest_rate"))
5. 关联计算转换值
通过from_curr、to_curr、year_month关联两个表,计算转换后的值:
df3 = df1_with_month.join(df2_latest_rate, on=["from_curr", "to_curr", "year_month"], how="left") \ .withColumn("converted_value", F.col("value_to_convert") * F.col("latest_rate"))
关键注意事项
- 确保
Date字段为PySpark的DateType,若原始数据是字符串,需用F.to_date()指定格式转换。 - 关联方式选择
left可保留所有待转换数据,若匹配不到汇率,converted_value会为null,可通过F.coalesce()设置默认值或标记。 - 若汇率表无重复记录,可跳过步骤4直接关联。
内容的提问来源于stack exchange,提问作者f.ivy
相关产品推荐
相关产品推荐

