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

如何用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截图:
sparkUIScreenshot

当前实现代码

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连接重置错误的配置调整

  1. 清理本地模式无效配置
    你当前使用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]。

  2. 优化网络与心跳参数
    延长超时时间,避免心跳中断:

    .config("spark.network.timeout", "1200s")
    .config("spark.executor.heartbeatInterval", "30s")
    
  3. 适配本地内存配置
    本地模式下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)的代码优化

  1. 减少Driver端内存占用
    避免在Driver端用列表累积大量数据,改用Pandas中转后再转Spark DataFrame:

    if len(trades) == 10000:
        pd_df = pd.DataFrame(trades)
        df = spark.createDataFrame(pd_df, schema=schema)
        # 后续写入逻辑不变
        trades = []
    
  2. 优化Parquet写入逻辑

    • 不要为每个批次创建单独目录,直接写入根目录,减少IO开销:
      robust_write(df, output_dir)  # 去掉batch_*子目录
      
    • 若exchange基数极小,暂时关闭partitionBy,避免生成过多小文件:
      df.write.mode("append").parquet(path)
      
  3. 显式释放资源
    每批写入后回收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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 21:47:11