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

PySpark实现多表匹配并将费率列值追加至保单表的需求

PySpark 费率匹配与结果追加实现方案

1. 数据结构确认与示例初始化

先明确两张核心表的结构,以下是示例数据初始化代码(可替换为你的实际数据源):

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, split, when, lit, element_at, cast

# 初始化Spark会话
spark = SparkSession.builder.appName("PolicyRateMatching").getOrCreate()

# 示例 rates_table['base'] 结构:cov, vers, rf, rf_start, rf_end, rf_cat, amt, coeff, rate
rates_data = [
    ("A", "V1", "AGE", 18, 30, "MALE", 100, 0.8, 0.05),
    ("A", "V1", "AGE", 31, 50, "MALE", 150, 0.9, 0.06),
    ("A", "V1", "OCC", None, None, "OFFICE", 200, 1.0, 0.07),
    ("B", "V2", "CAR", 0, 5, "SEDAN", 300, 0.7, 0.04)
]
rates_table = spark.createDataFrame(rates_data, ["cov", "vers", "rf", "rf_start", "rf_end", "rf_cat", "amt", "coeff", "rate"])

# 示例 policies_table 结构:pol_no, vers, cov, rf_1, rf_2, rl_1, rl_2
policies_data = [
    ("POL001", "V1", "A", "AGE,OCC", "CAR", "25,MALE", "3,SEDAN"),
    ("POL002", "V2", "B", "CAR", "AGE", "4,SUV", "35,FEMALE")
]
policies_table = spark.createDataFrame(policies_data, ["pol_no", "vers", "cov", "rf_1", "rf_2", "rl_1", "rl_2"])

2. 拆分rf/rl字段为数组

由于rf和rl字段是逗号分隔的复合值,先拆分为数组以便逐个处理匹配规则:

policies_split = policies_table \
    .withColumn("rf_1_arr", split(col("rf_1"), ",")) \
    .withColumn("rf_1_key", element_at(col("rf_1_arr"), 1))  # 提取rf_1首个值作为匹配键
    .withColumn("rl_1_arr", split(col("rl_1"), ",")) \
    .withColumn("rf_2_arr", split(col("rf_2"), ",")) \
    .withColumn("rf_2_key", element_at(col("rf_2_arr"), 1))  # 提取rf_2首个值作为匹配键
    .withColumn("rl_2_arr", split(col("rl_2"), ","))

3. 分组合规匹配逻辑

处理rf_1/rl_1组匹配

# 关联rates_table获取rf_1对应的费率规则
match_rf1 = policies_split.join(
    rates_table.alias("r1"),
    (col("cov") == col("r1.cov")) & 
    (col("vers") == col("r1.vers")) & 
    (col("rf_1_key") == col("r1.rf")),
    how="left"
)

# 执行匹配规则:区间判断+分类匹配
match_rf1 = match_rf1 \
    .withColumn("rl_1_first", element_at(col("rl_1_arr"), 1)) \
    .withColumn("rl_1_category", element_at(col("rl_1_arr"), 2)) \
    # rf_start非空时,判断rl首个值是否在区间内
    .withColumn("is_in_range", when(
        col("r1.rf_start").isNotNull(),
        cast(col("rl_1_first"), "int").between(col("r1.rf_start"), col("r1.rf_end"))
    ).otherwise(lit(True))) \
    # 分类值匹配逻辑:rf_start非空时用rl第二个值,否则用rl所有值
    .withColumn("is_cat_match", when(
        col("r1.rf_start").isNotNull(),
        col("rl_1_category") == col("r1.rf_cat")
    ).otherwise(
        col("rl_1_first") == col("r1.rf_cat")
    )) \
    # 筛选符合双条件的费率项
    .filter(col("is_in_range") & col("is_cat_match")) \
    # 重命名匹配结果字段
    .withColumnRenamed("r1.amt", "amt_1") \
    .withColumnRenamed("r1.coeff", "coeff_1") \
    .withColumnRenamed("r1.rate", "rate_1")

处理rf_2/rl_2组匹配

# 关联rates_table获取rf_2对应的费率规则
match_rf2 = match_rf1.join(
    rates_table.alias("r2"),
    (col("cov") == col("r2.cov")) & 
    (col("vers") == col("r2.vers")) & 
    (col("rf_2_key") == col("r2.rf")),
    how="left"
)

# 复用匹配规则逻辑
match_rf2 = match_rf2 \
    .withColumn("rl_2_first", element_at(col("rl_2_arr"), 1)) \
    .withColumn("rl_2_category", element_at(col("rl_2_arr"), 2)) \
    .withColumn("is_in_range_2", when(
        col("r2.rf_start").isNotNull(),
        cast(col("rl_2_first"), "int").between(col("r2.rf_start"), col("r2.rf_end"))
    ).otherwise(lit(True))) \
    .withColumn("is_cat_match_2", when(
        col("r2.rf_start").isNotNull(),
        col("rl_2_category") == col("r2.rf_cat")
    ).otherwise(
        col("rl_2_first") == col("r2.rf_cat")
    )) \
    .filter(col("is_in_range_2") & col("is_cat_match_2")) \
    .withColumnRenamed("r2.amt", "amt_2") \
    .withColumnRenamed("r2.coeff", "coeff_2") \
    .withColumnRenamed("r2.rate", "rate_2")

4. 生成最终输出表

保留原保单表的所有字段,追加匹配得到的费率字段:

final_table = match_rf2.select(
    "pol_no", "vers", "cov", "rf_1", "rf_2", "rl_1", "rl_2",
    "amt_1", "coeff_1", "rate_1", "amt_2", "coeff_2", "rate_2"
)

# 查看结果
final_table.show(truncate=False)

核心规则对应说明

  • 表关联:通过cov+vers字段关联,同时用policies表rf列的首个值匹配rates表的rf字段;
  • 区间匹配:仅当rates表rf_start非空时生效,将rl列首个值转为整数后判断是否在rf_start与rf_end区间内;
  • 分类匹配:rl列后续值(或rf_start为空时的所有rl值)与rates表rf_cat字段完全匹配;
  • 分组追加:匹配得到的amt/coeff/rate按rf_1/rl_1、rf_2/rl_2分组,分别命名为amt_1/coeff_1等追加到原表。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 02:17:13