如何用PySpark处理生成器数据并写入Parquet?配置与报错排查
问题
我有一个返回可迭代对象的数据生成器,用于获取指定日期范围内的数据,总量接近10亿条。目标是将数据写入本地文件系统,后续通过PySpark ReadStream读取并写入Cassandra(该环节已完成),本次问题聚焦数据获取与本地写入。
我尝试实现以下流程:
- 使用生成器获取数据;
- 累积批量数据;
- 当批量达到指定大小,创建Spark DataFrame;
- 将DataFrame写入Parquet格式文件。
但过程中遇到段错误(core dumped)和Java连接重置错误。作为PySpark新手,希望获取Spark配置建议,解决频繁出现的如下错误:
Failed to write to data/data/polygon/trades/batch_99 on attempt 1: An error occurred while calling o1955.parquet.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 8 in stage 99.0 failed 1 times, most recent failure: Lost task 8.0 in stage 99.0 (TID 3176) (furkan-desktop executor driver): java.net.SocketException: Connection reset
Spark UI截图:
当前实现代码
from datetime import datetime import logging import time from dotenv import load_dotenv import pandas as pd import os from pyspark.sql import SparkSession from pyspark.sql.types import ( StructType, StructField, IntegerType, StringType, LongType, DoubleType, ArrayType, ) import uuid from polygon import RESTClient logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s", filename="spark_logs/logfile.log", filemode="w", ) from_date = datetime(2021, 3, 10) to_date = datetime(2021, 3, 31) load_dotenv() client = RESTClient(os.getenv("POLYGON_API_KEY")) # Create Spark session spark = ( SparkSession.builder.appName("TradeDataProcessing") .master("local[*]") .config("spark.driver.memory", "16g") .config("spark.executor.instances", "8") .config("spark.executor.memory", "16g") .config("spark.executor.memoryOverhead", "4g") .config("spark.executor.cores", "4") .config("spark.memory.offHeap.enabled", "true") .config("spark.memory.offHeap.size", "4g") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .config("spark.kryoserializer.buffer.max", "512m") .config("spark.network.timeout", "800s") .config("spark.executor.heartbeatInterval", "20000ms") .config("spark.dynamicAllocation.enabled", "true") .config("spark.dynamicAllocation.minExecutors", "1") .config("spark.dynamicAllocation.maxExecutors", "8") .config("spark.dynamicAllocation.initialExecutors", "4") .getOrCreate() ) # Define the schema corresponding to the JSON structure schema = StructType( [ StructField("exchange", IntegerType(), False), StructField("id", StringType(), False), StructField("participant_timestamp", LongType(), False), StructField("price", DoubleType(), False), StructField("size", DoubleType(), False), StructField("conditions", ArrayType(IntegerType()), True), ] ) def ensure_directory_exists(path): """Ensure directory exists, create if it doesn't""" if not os.path.exists(path): os.makedirs(path) # Convert dates to timestamps or use them directly based on your API requirements from_timestamp = from_date.timestamp() * 1e9 # Adjusting for nanoseconds to_timestamp = to_date.timestamp() * 1e9 # Initialize the trades iterator with the specified parameters trades_iterator = client.list_trades( "X:BTC-USD", timestamp_gte=int(from_timestamp), timestamp_lte=int(to_timestamp), limit=1_000, sort="asc", order="asc", ) trades = [] file_index = 0 output_dir = "data/data/polygon/trades" # Output directory ensure_directory_exists(output_dir) # Make sure the output directory exists def robust_write(df, path, max_retries=3, retry_delay=5): """Attempts to write a DataFrame to a path with retries on failure.""" for attempt in range(max_retries): try: df.write.partitionBy("exchange").mode("append").parquet(path) print(f"Successfully written to {path}") return except Exception as e: logging.error(f"Failed to write to {path} on attempt {attempt+1}: {e}") time.sleep(retry_delay) # Wait before retrying logging.critical(f"Failed to write to {path} after {max_retries} attempts.") for trade in trades_iterator: trade_data = { "exchange": int(trade.exchange), "id": str(uuid.uuid4()), "participant_timestamp": trade.participant_timestamp, "price": float(trade.price), "size": float(trade.size), "conditions": trade.conditions if trade.conditions else [], } trades.append(trade_data) if len(trades) == 10000: df = spark.createDataFrame(trades, schema=schema) file_name = f"{output_dir}/batch_{file_index}" robust_write(df, file_name) trades = [] file_index += 1 if trades: df = spark.createDataFrame(trades, schema=schema) file_name = f"{output_dir}/batch_{file_index}" robust_write(df, file_name)
解决方案与配置优化建议
一、针对Java连接重置错误的配置调整
清理本地模式无效配置
你当前使用master("local[*]")本地单JVM模式,spark.executor.instances、spark.dynamicAllocation.*这类集群参数完全无效,会导致资源管理混乱。删除以下配置:.config("spark.executor.instances", "8") .config("spark.dynamicAllocation.enabled", "true") .config("spark.dynamicAllocation.minExecutors", "1") .config("spark.dynamicAllocation.maxExecutors", "8") .config("spark.dynamicAllocation.initialExecutors", "4")本地模式下直接用
local[*]自动使用所有CPU核心,或指定核心数如local[8]。优化网络与心跳参数
延长超时时间,避免心跳中断:.config("spark.network.timeout", "1200s") .config("spark.executor.heartbeatInterval", "30s")适配本地内存配置
本地模式下Driver内存即为JVM总内存,调整内存参数适配本地机器:.config("spark.driver.memory", "24g") # 根据实际机器内存调整,如32G内存设24G .config("spark.driver.memoryOverhead", "4g") .config("spark.memory.offHeap.enabled", "true") .config("spark.memory.offHeap.size", "8g")
二、解决段错误(core dumped)的代码优化
减少Driver端内存占用
避免在Driver端用列表累积大量数据,改用Pandas中转后再转Spark DataFrame:if len(trades) == 10000: pd_df = pd.DataFrame(trades) df = spark.createDataFrame(pd_df, schema=schema) # 后续写入逻辑不变 trades = []优化Parquet写入逻辑
- 不要为每个批次创建单独目录,直接写入根目录,减少IO开销:
robust_write(df, output_dir) # 去掉batch_*子目录 - 若
exchange基数极小,暂时关闭partitionBy,避免生成过多小文件:df.write.mode("append").parquet(path)
- 不要为每个批次创建单独目录,直接写入根目录,减少IO开销:
显式释放资源
每批写入后回收DataFrame资源,避免内存泄漏:robust_write(df, file_name) df.unpersist() del df
三、其他实用建议
- 增大批量大小:将每批数据量从10000调整到50000或100000,减少Spark作业提交次数。
- 优化日志配置:添加日志参数定位详细错误:
.config("spark.driver.extraJavaOptions", "-Dlog4j.configuration=file:log4j.properties") - 检查磁盘状态:确保写入目录所在磁盘有足够空间且IO性能达标,避免磁盘瓶颈导致超时。
内容的提问来源于stack exchange,提问作者Furkan Öztürk

