如何在PySpark中实现NearMiss等智能欠采样技术用于欺诈检测
基于PySpark实现智能欠采样处理类别不平衡问题
推荐的PySpark不平衡数据处理库
- pyspark-imbalanced-learning:该库基于原生Spark实现了包括TomekLinks、NearMiss、ClusterCentroids、ENN在内的多种智能欠采样/过采样算法,无需将Spark DataFrame转换为本地Pandas,完美适配PySpark生态。安装命令:
pip install pyspark-imbalanced-learning
原生Spark实现TomekLinks欠采样
TomekLinks的核心逻辑是识别并剔除多数类中与少数类样本形成“Tomek对”的样本(即同类最近邻距离小于异类最近邻距离的样本),以下是原生Spark的实现代码:
from pyspark.sql import SparkSession from pyspark.sql.functions import col from pyspark.ml.feature import VectorAssembler, NearestNeighbors # 初始化Spark会话 spark = SparkSession.builder.appName("TomekLinksUndersampling").getOrCreate() # 加载数据集(替换为你的实际数据路径) df = spark.read.csv("fraudTrain.csv", header=True, inferSchema=True) # 分离少数类(欺诈样本)和多数类(正常样本) minority_df = df.filter(col("is_fraud") == 1) majority_df = df.filter(col("is_fraud") == 0) # 定义特征列(排除目标列和非特征字段) feature_cols = [c for c in df.columns if c not in ["is_fraud", "trans_date_trans_time", "cc_num", "merchant"]] assembler = VectorAssembler(inputCols=feature_cols, outputCol="features") # 转换为特征向量格式 minority_vec = assembler.transform(minority_df).select("features", "is_fraud") majority_vec = assembler.transform(majority_df).select("features", "is_fraud") # 计算每个多数类样本到最近少数类样本的距离 nn_minority = NearestNeighbors(k=1, inputCol="features", outputCol="nearest_minority") model_minority = nn_minority.fit(minority_vec) majority_with_min_dist = model_minority.transform(majority_vec).select( "features", "is_fraud", col("nearest_minority.distance").alias("min_minority_dist") ) # 计算每个多数类样本到最近同类样本的距离(k=2排除自身) nn_majority = NearestNeighbors(k=2, inputCol="features", outputCol="nearest_majority") model_majority = nn_majority.fit(majority_vec) majority_with_self_dist = model_majority.transform(majority_vec).select( "features", col("nearest_majority")[1]["distance"].alias("min_majority_dist") ) # 筛选非Tomek对的多数类样本 majority_filtered = majority_with_min_dist.join( majority_with_self_dist, on="features", how="inner" ).filter(col("min_majority_dist") >= col("min_minority_dist")).select( majority_with_min_dist["features"], majority_with_min_dist["is_fraud"] ) # 合并少数类与筛选后的多数类,得到平衡数据集 balanced_df = minority_vec.union(majority_filtered) # 可选:转换回原始特征格式(剔除特征向量列) balanced_df = balanced_df.join(df, on="features", how="inner").drop("features")
PySpark Pandas适配方案
若偏好使用PySpark Pandas,可直接调用pyspark-imbalanced-learning的适配接口,避免本地Pandas的兼容性问题:
import pyspark.pandas as ps from pyspark_imbalanced_learning.under_sampling import TomekLinks # 加载数据为PySpark Pandas DataFrame ps_df = ps.read_csv("fraudTrain.csv") # 分离特征与目标变量 X = ps_df.drop(["is_fraud", "trans_date_trans_time", "cc_num", "merchant"], axis=1) y = ps_df["is_fraud"] # 执行TomekLinks欠采样 tl_sampler = TomekLinks() X_resampled, y_resampled = tl_sampler.fit_resample(X, y) # 合并为平衡数据集 balanced_ps_df = ps.concat([X_resampled, y_resampled], axis=1)
内容的提问来源于stack exchange,提问作者Sumit
相关产品推荐
相关产品推荐

