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

