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。
实现步骤
- 过滤ODM映射表,仅保留
odm_id='abcefgh'且odm_relation='EQ'的有效等值映射记录 - 将交易表与过滤后的映射表通过
tr_type = input_value关联,获取每个交易对应的return_value - 使用
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
相关产品推荐
相关产品推荐

