如何在Azure Databricks中从指定文件夹(2022年起)读取流数据
解决Azure Databricks流读取指定年份数据及内存不足问题
一、修复modifiedAfter无效问题:改用路径过滤
你的数据按yyyymmdd格式的文件夹分层存储,直接通过路径模式匹配2022年及以后的文件夹,比依赖文件修改时间更可靠——modifiedAfter失效通常是因为文件实际修改时间与文件夹命名时间不匹配,或是路径扫描逻辑未正确识别该参数。
修改输入路径为匹配2022年及以后的模式:
# 匹配2022-2029年的所有子文件夹及文件(适配yyyymm或yyyymmdd层级) inputPath = '/mnt/ASN-1.0/202[2-9]*/**' # 若你的层级是yyyy/mm/dd,可更精确匹配: # inputPath = '/mnt/ASN-1.0/202[2-9][0-9]/**'
二、解决全量读取内存不足问题:多维度优化配置
针对25GB数据的内存压力,从流读取控制、Spark资源配置两方面优化:
1. 限制单次触发的文件读取量
避免一次性加载过多文件到内存,通过maxFilesPerTrigger控制每次流触发处理的文件数:
.option("maxFilesPerTrigger", 100) # 根据集群规模调整,比如50-200之间
2. 优化Spark executor资源与并行度
在代码中添加以下配置,提升内存使用效率、分散计算压力:
# 调整executor内存与核心数(根据你的集群规格,示例为4核16G配置) spark.conf.set("spark.executor.memory", "16g") spark.conf.set("spark.executor.cores", 4) # 调整并行度,适配executor数量 spark.conf.set("spark.sql.shuffle.partitions", 200) spark.conf.set("spark.default.parallelism", 200) # 开启堆外内存缓解堆内压力 spark.conf.set("spark.memory.offHeap.enabled", "true") spark.conf.set("spark.memory.offHeap.size", "8g") # 禁用大表自动广播,避免内存溢出 spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
3. 依托Checkpoint实现增量读取
你使用的trigger(once=True)配合Checkpoint,第一次运行会读取2022年及以后的历史数据,后续运行时Checkpoint会自动跳过已处理文件,仅读取新增数据,无需重复扫描全量文件。注意Checkpoint路径必须唯一,不可与其他流任务共用。
三、完整优化代码示例
import json from pyspark.sql.types import StructType # 基础配置 checkpointPath = '/mnt/checkpointasnmod1' schemapath = '/FileStore/tables/scema-1.txt' inputPath = '/mnt/ASN-1.0/202[2-9]*/**' outputPath = '/mnt/rawoutputpartially1' # 加载Schema schemaJson = spark.read.text(schemapath).first()[0] schemaStruct = StructType.fromJson(json.loads(schemaJson)) # Spark资源优化配置 spark.conf.set("spark.executor.memory", "16g") spark.conf.set("spark.executor.cores", 4) spark.conf.set("spark.sql.shuffle.partitions", 200) spark.conf.set("spark.memory.offHeap.enabled", "true") spark.conf.set("spark.memory.offHeap.size", "8g") spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1) # 流处理逻辑 df = spark.readStream.format("cloudFiles") .option("cloudFiles.useNotifications","false") .option("cloudFiles.validateOptions","false") .option("cloudFiles.format","text") .option("wholetext","true") .option("maxFilesPerTrigger", 100) .load(inputPath, schema=None) .writeStream .trigger(once=True) .outputMode("append") .queryName("ASN_Stream_Processing") .foreachBatch(transformasn) .option("checkpointLocation", checkpointPath) .start() df.awaitTermination()
内容的提问来源于stack exchange,提问作者anuj
相关产品推荐
相关产品推荐

