如何查看PySpark DataFrame单条记录大小?排查Synapse行超限错误
定位PySpark DataFrame中超大记录的可行方法
针对Synapse PolyBase的1MB行大小限制问题,以下几种方法可以帮你定位超出限制的记录:
1. 序列化计算每行精准大小
通过UDF将每行数据序列化,直接计算字节数,这个结果和实际传输时的大小接近:
from pyspark.sql import functions as F from pyspark.sql.types import LongType import pickle def get_row_size(row): return len(pickle.dumps(row)) row_size_udf = F.udf(get_row_size, LongType()) # 给DataFrame添加行大小列,筛选超1MB的记录 df_with_size = df.withColumn("row_size_bytes", row_size_udf(F.struct(*df.columns))) large_rows = df_with_size.filter(F.col("row_size_bytes") > 1000000) # 查看问题记录,truncate=False避免截断内容 large_rows.show(truncate=False)
2. 按列累加估算行大小
如果想知道具体是哪些列占空间,可以逐个计算列的字节数再累加:
from pyspark.sql import functions as F column_size_exprs = [] for col_name, col_type in df.dtypes: if col_type.startswith("string"): # octet_length直接获取字符串的字节数(UTF-8编码) column_size_exprs.append(F.octet_length(F.col(col_name)).alias(f"{col_name}_bytes")) elif col_type in ["int", "integer"]: column_size_exprs.append(F.lit(4).alias(f"{col_name}_bytes")) elif col_type in ["bigint", "long"]: column_size_exprs.append(F.lit(8).alias(f"{col_name}_bytes")) elif col_type == "double": column_size_exprs.append(F.lit(8).alias(f"{col_name}_bytes")) elif col_type == "float": column_size_exprs.append(F.lit(4).alias(f"{col_name}_bytes")) elif col_type.startswith("array") or col_type.startswith("struct"): # 复杂类型转JSON后算字节数 column_size_exprs.append(F.octet_length(F.to_json(F.col(col_name))).alias(f"{col_name}_bytes")) else: # 其他类型按实际情况调整默认值 column_size_exprs.append(F.lit(16).alias(f"{col_name}_bytes")) # 计算每行总大小并筛选 df_with_sizes = df.select(*df.columns, *column_size_exprs) df_with_total = df_with_sizes.withColumn( "total_row_size", sum(F.col(f"{col}_bytes") for col, _ in df.dtypes) ) large_rows = df_with_total.filter(F.col("total_row_size") > 1000000) large_rows.show(truncate=False)
3. 先定位问题Parquet文件
如果Parquet文件数量多,可以先通过元数据找到包含超大行的文件,缩小排查范围:
from pyspark.sql import functions as F # 获取所有Parquet文件路径 parquet_paths = df.inputFiles() for path in parquet_paths: # 统计每个文件的最大行大小 file_stats = spark.read.parquet(path).select( F.lit(path).alias("file_path"), F.max(F.octet_length(F.to_json(F.struct(*df.columns)))).alias("max_row_size") ) file_stats.show()
找到max_row_size超过1MB的文件后,再针对该文件进行逐行排查即可。
内容的提问来源于stack exchange,提问作者LearneR
相关产品推荐
相关产品推荐

