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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 17:42:43