如何在PySpark中实现计数模型的负二项回归
PySpark 3.0+ 版本的pyspark.ml模块已原生支持负二项分布的广义线性回归,可直接实现你需要的计数模型负二项回归,对应实现方案如下:
实现步骤与代码
1. 依赖导入
from pyspark.sql import SparkSession from pyspark.ml.regression import GeneralizedLinearRegression from pyspark.ml.feature import RFormula from pyspark.ml import Pipeline import pyspark.sql.functions as F
2. 数据预处理
# 初始化Spark会话 spark = SparkSession.builder.appName("NegativeBinomialRegression").getOrCreate() # 加载合约数据,替换为你自己的数据源加载逻辑 contracts_df = spark.read.csv("你的数据路径", header=True, inferSchema=True) # 生成offset列,对应原有逻辑中的offset(log(Exposure)) contracts_df = contracts_df.withColumn("log_exposure", F.log(F.col("Exposure")))
3. 特征处理与模型定义
这里用RFormula和你原有R、statsmodels的公式写法完全对齐,无需手动处理分类变量编码:
# 定义特征、标签映射规则 formula = RFormula( formula="ClaimNb ~ DriverAge + CarAge + Density + Brand + Power + Gas", featuresCol="features", labelCol="ClaimNb" ) # 定义负二项回归模型 nb_glm = GeneralizedLinearRegression( family="negative.binomial", link="log", offsetCol="log_exposure", # 如需固定分散参数和statsmodels的alpha=1对齐,取消下一行注释即可 # variancePower=1.0 )
4. 模型训练与结果查看
# 构造流水线执行特征处理+模型训练 pipeline = Pipeline(stages=[formula, nb_glm]) trained_pipeline = pipeline.fit(contracts_df) # 提取训练完成的负二项回归模型 nb_model = trained_pipeline.stages[-1] # 输出模型统计结果,对应原有summary功能 print("模型系数:", nb_model.coefficients) print("模型截距:", nb_model.intercept) print("模型完整摘要:\n", nb_model.summary) # 可单独提取需要的指标 print("AIC:", nb_model.summary.aic) print("系数P值:", nb_model.summary.pValues) print("残差自由度:", nb_model.summary.residualDegreeOfFreedom)
注意事项
- 如果你的PySpark版本低于3.0,没有内置负二项分布支持,优先升级到Spark 3.x版本即可获得原生支持。
variancePower参数对应负二项分布的分散参数,和statsmodels中NegativeBinomial(alpha=xxx)的alpha参数含义完全一致,不指定该参数时模型会自动拟合最优的分散参数。
内容的提问来源于stack exchange,提问作者Ramsai Jagadish
相关产品推荐
相关产品推荐

