使用Delta Live Tables运行代码时遇Photon内存不足错误求助
Delta Live Tables 内存溢出问题排查与解决
运行代码
from pyspark.sql.types import StructType, StructField, StringType from pyspark.sql.functions import col import dlt @dlt.view def data(): return (spark.readStream .option("skipChangeCommits", "true") .format("delta") .table("table") .withColumnRenamed("col","col_name") ) schema = StructType([ StructField("col_name", StringType(), True), ]) dlt.create_streaming_table( name="table", spark_conf={}, table_properties={"quality":"bronze"}, partition_cols=["col"], schema=schema ) dlt.apply_changes( target="table", source="data", keys=["userId"], sequence_by=col("col_name"), stored_as_scd_type=1 )
报错信息
原因:org.apache.spark.memory.SparkOutOfMemoryError: Photon执行查询时内存不足。在任务中的FileWriterNode(id=34323, output_schema=[])内的ParquetDictionaryEncoder中,Photon无法为哈希表变长键数据预留349.6 MiB内存。
问题分析与解决
代码逻辑错误(加剧内存问题)
- 循环依赖:
data视图读取目标表table,而apply_changes又将data的数据写入table,形成循环处理,会导致数据无限循环处理,直接引发内存过载。 - 字段不匹配:
create_streaming_table指定的分区列col不存在于定义的schema中(schema只有col_name),会导致分区逻辑异常。apply_changes的keys=["userId"],但schema中无userId字段;sequence_by使用字符串类型的col_name,不符合序列列需为有序类型(如时间戳、递增ID)的要求,会导致数据处理逻辑混乱,增加内存消耗。
内存溢出解决措施
调整Photon内存配置:
- 增加Executor内存:设置
spark.executor.memory为更大值(如16g),同时调整spark.executor.cores匹配内存资源。 - 提高哈希表内存上限:设置
spark.sql.photon.memory.hashTable.maxMemory为更高值(如512m),允许Photon为哈希表分配更多内存。 - 临时关闭Photon:若不需要Photon加速,可设置
spark.sql.execution.phototon.enabled=false,切换回Spark原生执行引擎。
- 增加Executor内存:设置
修复代码逻辑:
- 解除循环依赖:修改
data视图的数据源为实际业务表,而非目标表table。 - 修正字段匹配:
- 将
partition_cols改为schema中存在的字段(如col_name),或调整schema包含col字段。 keys替换为实际存在的主键字段,sequence_by改用时间戳或递增ID类有序字段。
- 将
- 优化流处理触发间隔:在
readStream中添加trigger,减少每个微批处理的数据量,例如:spark.readStream.trigger(processingTime="5 minutes")
- 解除循环依赖:修改
优化数据存储:
- 避免高基数分区:选择低基数字段作为分区列,减少小文件数量。
- 调整文件大小:设置
spark.sql.files.maxPartitionBytes(默认128M)为合适值,控制输出文件大小,降低内存处理压力。
内容的提问来源于stack exchange,提问作者Panda
相关产品推荐
相关产品推荐

