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

PySpark中基于列级ODM映射实现交易金额正负调整

PySpark实现基于ODM映射的交易金额取反逻辑

现有两个PySpark DataFrame:

  • 交易数据表df1:
+---+-------+--------+
|id |tr_type|nominal |
+---+-------+--------+
|1  |K      |2.0     |
|2  |ZW     |7.0     |
|3  |V      |12.5    |
|4  |VW     |9.0     |
|5  |CI     |5.0     |
+---+-------+--------+
  • ODM映射表(仅关注odm_id为abcefgh的记录):
+-------+------------+------------+-----------+
|odm_id |return_value|odm_relation|input_value|
+-------+------------+------------+-----------+
|abcefgh|B           |EQ          |K          |
|abcefgh|B           |EQ          |ZW         |
|abcefgh|S           |EQ          |V          |
|abcefgh|S           |EQ          |VW         |
|abcefgh|I           |EQ          |CI         |
+-------+------------+------------+-----------+

需求:当交易类型tr_type通过ODM映射得到的return_value为'S'(卖出交易)时,将nominal字段取反,生成新字段nominal_new。


实现步骤

  1. 过滤ODM映射表,仅保留odm_id='abcefgh'且odm_relation='EQ'的有效等值映射记录
  2. 将交易表与过滤后的映射表通过tr_type = input_value关联,获取每个交易对应的return_value
  3. 使用when...otherwise条件逻辑生成nominal_new:若return_value='S'则取反nominal,否则保留原值

完整代码示例

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when

# 初始化SparkSession
spark = SparkSession.builder.appName("ODM_Nominal_Adjustment").getOrCreate()

# 创建交易数据表df1
data_df1 = [
    (1, "K", 2.0),
    (2, "ZW", 7.0),
    (3, "V", 12.5),
    (4, "VW", 9.0),
    (5, "CI", 5.0)
]
df1 = spark.createDataFrame(data_df1, ["id", "tr_type", "nominal"])

# 创建ODM映射表
data_odm = [
    ("abcefgh", "B", "EQ", "K"),
    ("abcefgh", "B", "EQ", "ZW"),
    ("abcefgh", "S", "EQ", "V"),
    ("abcefgh", "S", "EQ", "VW"),
    ("abcefgh", "I", "EQ", "CI")
]
odm_df = spark.createDataFrame(data_odm, ["odm_id", "return_value", "odm_relation", "input_value"])

# 过滤有效映射
filtered_odm = odm_df.filter((col("odm_id") == "abcefgh") & (col("odm_relation") == "EQ"))

# 关联并计算nominal_new
result_df = df1.join(
    filtered_odm,
    df1.tr_type == filtered_odm.input_value,
    "left"
).select(
    df1.id,
    df1.tr_type,
    df1.nominal,
    when(col("return_value") == "S", -col("nominal")).otherwise(col("nominal")).alias("nominal_new")
)

# 展示结果
result_df.show()

输出结果

+---+-------+-------+-----------+
|id |tr_type|nominal|nominal_new|
+---+-------+-------+-----------+
|1  |K      |2.0    |2.0        |
|2  |ZW     |7.0    |7.0        |
|3  |V      |12.5   |-12.5      |
|4  |VW     |9.0    |-9.0       |
|5  |CI     |5.0    |5.0        |
+---+-------+-------+-----------+

关键说明

  • 使用left join确保所有交易记录都被保留,即使出现未在映射表中的tr_type(此时return_value为null,nominal_new会保留原值)
  • 直接用内置when函数完成条件计算,无需自定义UDF,性能更优

内容的提问来源于stack exchange,提问作者santhosh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 22:45:37