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

PySpark pandas DataFrame apply仅处理前1000行问题求助

问题分析与解决方案

核心问题1:pyspark.pandas.DataFrame.apply仅处理1000行的原因

你遇到的现象源于两个关键问题:

  • pyspark.pandas的测试机制:默认情况下apply会先运行1000行的测试批次,验证自定义函数的合法性,这个测试在Driver端执行。
  • 全局变量在分布式环境无效:你通过dict_list.append()修改全局列表的方式完全不适配Spark分布式架构——每个Executor都有独立内存空间,它们的dict_list副本不会同步回Driver端,只有测试批次的1000行数据被Driver端处理并写入了你的列表。

核心问题2:Spark DataFrame无法识别数字列名

Spark支持数字列名,只需用**反引号(`)**包裹引用,或者提前将列名重命名为字符串,即可正常使用。


正确解决方案

方案1:使用普通Spark DataFrame(推荐,性能更优)

先解决列名问题,再用Spark UDF处理数据:

import ast
import os
import pandas as pd
import findspark
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import MapType, StringType

# 初始化Spark
findspark.init()
memory = "100g"
pyspark_submit_args = f" --driver-memory {memory} pyspark-shell"
os.environ["KMP_DUPLICATE_LIB_OK"] = "True"
os.environ["PYSPARK_SUBMIT_ARGS"] = pyspark_submit_args
os.environ["PYARROW_IGNORE_TIMEZONE"] = "1"

# 构建SparkSession
spark = SparkSession.builder \
    .config("spark.driver.maxResultSize", "500G") \
    .config("spark.executor.memory", "200G") \
    .config("spark.app.name", "Spark Updated Conf") \
    .config("spark.executor.cores", "5") \
    .config("spark.executor.instances", "5") \
    .config("spark.cores.max", "8") \
    .config("spark.driver.memory", "500G") \
    .config("spark.sql.shuffle.partitions", "500") \
    .config("spark.sql.execution.arrow.enabled", "true") \
    .config("spark.default.parallelism", "500") \
    .config("spark.sql.execution.arrow.pyspark.enabled", "true") \
    .getOrCreate()

# 读取数据并重命名数字列名
df_pd = pd.read_parquet(
    "training_data_fv_unlabeled_spark_libsvm.parquet", engine="fastparquet"
)
# 假设列名是数字0,重命名为字符串列名
df_pd = df_pd.rename(columns={0: "feature_vector"})

# 转成Spark DataFrame
df_spark = spark.createDataFrame(df_pd)

# 定义处理函数和对应的UDF
def _parse_dict_key_to_str(dict_obj):
    # 补充你的_parse_dict_key_to_str实现,示例:将key转为字符串
    return {str(k): v for k, v in dict_obj.items()}

def fv_to_dict(fv_str):
    fv_str = "{" + fv_str.replace(" ", ", ") + "}"
    dict_obj = ast.literal_eval(fv_str)
    return _parse_dict_key_to_str(dict_obj)

# 注册UDF,指定返回类型(根据实际需求调整)
fv_udf = udf(fv_to_dict, MapType(StringType(), StringType()))

# 处理数据
df_processed = df_spark.withColumn("processed_dict", fv_udf("feature_vector"))

# 若需收集结果到Driver(注意:600万行可能内存压力大,建议优先保存到文件)
dict_list = df_processed.select("processed_dict").rdd.flatMap(lambda x: [x[0]]).collect()

# 或直接保存结果到Parquet
df_processed.write.parquet("processed_data.parquet")

方案2:使用pyspark.pandas(保持原有习惯,但修正逻辑)

若坚持用pyspark.pandas,不要依赖全局变量,而是让apply返回结果并触发计算:

# (前面的Spark初始化和数据读取重命名步骤与方案1一致)
import pyspark.pandas as ps

psdf = ps.DataFrame(df_pd)

# 修改函数,返回处理后的结果而非append到全局列表
def fv_to_dict(row):
    fv_str = row.feature_vector
    fv_str = "{" + fv_str.replace(" ", ", ") + "}"
    dict_obj = ast.literal_eval(fv_str)
    return _parse_dict_key_to_str(dict_obj)

# 执行apply并触发计算
psdf["processed_dict"] = psdf.apply(fv_to_dict, axis=1)

# 收集结果到Driver(注意内存压力)
dict_list = psdf["processed_dict"].to_list()

关键注意事项

  • 避免全局变量:Spark是分布式架构,Executor无法访问Driver的全局变量,所有数据处理逻辑都应通过返回值传递。
  • 数字列名处理:要么重命名列,要么在Spark中用反引号引用(比如df_spark.select(0))。
  • 内存压力:600万行数据直接collect()到Driver可能导致内存溢出,建议优先将结果写入分布式存储(如Parquet),而非收集到本地。

内容的提问来源于stack exchange,提问作者Saachi Shenoy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 17:35:23