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

极简Spark Streaming(读文件流写Delta)应用内存泄漏求助

Spark Streaming Delta Upsert 内存泄漏排查分析

问题描述

开发了一款极简Spark Streaming应用,功能为从S3读取小体积Parquet文件(KB级,每2~3小时生成一个),通过Upsert方式写入Delta表(表大小约190MB,包含约70万行数据)。应用部署在Spark+Kubernetes集群中,出现容器与Pod内存随时间持续增长,直至达到资源限制导致Pod崩溃的内存泄漏问题。

应用代码

import boto3
import os

from pathlib import Path

from delta import *
from delta.tables import *
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *


app_name = "my_app"
spark = SparkSession.builder.appName(app_name).config("spark.databricks.io.cache.enabled", "true").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
env = os.environ["ENVIRONMENT"]

landing_bucket = f"s3a://awesommmme-data-landing-zone-{env}"  # data_real/landing
cleansing_bucket = f"s3a://awesommmme-data-cleansing-zone-{env}"
checkpoint_bucket = f"s3a://awesommmme-data-spark-checkpoints-{env}"

landing_data_label_path = os.path.join(landing_bucket, "contract_label", "data")
cleansing_data_label_path = os.path.join(cleansing_bucket, "contract_label", "data")

_label_logic_schema = [
    StructField("token_address", StringType()),
    StructField("token_type", StringType()),
]
_raw_data_schema = _label_logic_schema + [
    StructField("update_dt", TimestampType()),
    StructField("year", StringType()),
    StructField("month", StringType()),
    StructField("day", StringType()),
    StructField("hour", StringType()),
    StructField("minute", StringType()),
]
raw_data_schema = StructType(_raw_data_schema)


deltaTable = None
def upsert_to_cleaned_delta(pdf, batchId):
    global deltaTable

    if deltaTable is None:
        try:
            deltaTable = DeltaTable.forPath(spark, cleansing_data_label_path)
        except Exception as e:
            import logging

            logging.error(str(e))
            pdf.write.format("delta").partitionBy("year", "month", "day", "hour", "minute").save(
                cleansing_data_label_path
            )
            deltaTable = DeltaTable.forPath(spark, cleansing_data_label_path)
            return

    deltaTable.alias("old_data").merge(
        pdf.alias("new_data"), "old_data.token_address = new_data.token_address"
    ).whenMatchedUpdateAll(
    ).whenNotMatchedInsertAll().execute()

 
def main():     
    _ = (
        spark.readStream.schema(raw_data_schema)
        .parquet(landing_data_label_path)
        .writeStream.format("delta")
        .option("maxFilesPerTrigger", 10)
        .option("checkpointLocation", "MY_CHECKPOINT_PATH")
        .foreachBatch(lambda pdf, batch_id: upsert_to_cleaned_delta(pdf, batch_id))
        .trigger(processingTime="1 second")
        .outputMode("append")
        .start()
    )
    spark.streams.awaitAnyTermination()


if __name__ == "__main__":
    main()

部署配置(values.yaml)

app: data-pipeline
applicationId: myapp
name: myapp-data
namespace: "{{ .Release.Namespace }}"

environment: dev


sparkConf:
  spark.executor.heartbeatInterval: "600s"
  spark.network.timeout: "3600s"

deps:
  jars:
    - delta-core_2.12-2.0.1.jar

restartPolicy:
  type: Always


driver:
  annotations: {}
  coreRequest: 100m
  coreLimit: 500m
  memory: 3g
executor:
  coreRequest: 100m
  coreLimit: 500m
  instances: 1
  memory: 2g

排查方向与优化建议

  • 全局变量deltaTable的资源持有问题:在foreachBatch中使用全局变量持有DeltaTable实例,长期持有可能导致元数据或未释放资源累积。建议每次batch时重新获取DeltaTable实例,或确保实例能被JVM正确垃圾回收。
  • 过短的触发间隔:设置processingTime="1 second"触发,但实际数据生成频率为2~3小时一次,频繁空触发会导致Driver累积调度对象与元数据。建议调整为trigger(availableNow=True)(Spark 3.3+版本支持,仅处理新到达数据后休眠),或设置更长的触发间隔(如10分钟)。
  • 非适配的缓存配置:spark.databricks.io.cache.enabled=true是Databricks环境专属优化,在开源Spark+K8s环境可能导致缓存无法正确清理,尝试关闭该配置观察内存变化。
  • JVM GC调优:当前未配置GC参数,默认GC可能无法及时回收对象。添加以下Spark配置开启G1GC并输出GC日志,分析内存回收情况:
    sparkConf:
      spark.executor.extraJavaOptions: "-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+PrintGCDetails -XX:+PrintGCTimeStamps"
      spark.driver.extraJavaOptions: "-XX:+UseG1GC -XX:MaxGCPauseMillis=200 -XX:+PrintGCDetails -XX:+PrintGCTimeStamps"
    
  • Checkpoint路径配置错误:代码中checkpointLocation硬编码为"MY_CHECKPOINT_PATH",需确认是否指向checkpoint_bucket下的有效路径。路径错误会导致Spark不断尝试创建无效checkpoint对象,累积内存占用。
  • Delta Merge操作的重复执行:每秒一次的触发会导致Driver重复执行Merge逻辑,即使没有新数据。结合触发间隔调整,减少不必要的Merge操作,降低内存累积。
  • 堆内存分析:通过kubectl top pod监控Pod内存,使用jmap -dump:format=b,file=heap.hprof <pid>生成堆转储文件,分析堆内存中对象分布,定位内存占用大户。
  • 版本兼容性验证:确认delta-core_2.12-2.0.1与Spark版本兼容(Delta 2.0.1对应Spark 3.2.x),版本不兼容可能导致内存泄漏或不稳定。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 11:20:25