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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 10:21:12