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
相关产品推荐
相关产品推荐

