如何使Synapse Notebook仅处理Azure存储账户中的新增数据?
高效识别并处理Azure存储中新增数据的方案
方法一:通过文件最后修改时间筛选
- 在PySpark Notebook中,读取存储容器内的文件元数据(包含最后修改时间),可以用
spark.read.format("binaryFile").load("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/"),返回的DataFrame会自带path和modificationTime字段 - 在湖数据库中维护一张
processed_metadata表,只存一条记录,用来记录上次处理的最大时间戳 - 每次运行Notebook时,先读取该时间戳,过滤出
modificationTime大于此时间戳的文件进行处理 - 处理完成后,更新
processed_metadata表的时间戳为本次处理的最大modificationTime
方法二:利用文件路径命名规则匹配
- 你的文件路径遵循
YYYY/WeekN/WeekN.xlsx的固定格式,可在湖数据库中建一张processed_weeks表,记录已处理的年份-周数组合(比如2022-Week1) - 运行Notebook时,先列出存储容器内所有文件的路径,通过字符串拆分提取每个文件对应的年份和周数
- 将提取出的
年份-周数与processed_weeks表做左连接,筛选出未出现在表中的条目,定位对应文件处理 - 处理完成后,将本次处理的
年份-周数插入到processed_weeks表中
方法三:管道参数化+定时触发
- 因为数据是每周新增,直接在Synapse管道中设置每周定时触发器
- 在管道中添加
current_year和current_week参数,通过触发时间自动计算当前年份和周数(用表达式@formatDateTime(utcNow(), 'yyyy')和@formatDateTime(utcNow(), 'ww')) - 将参数传递给PySpark Notebook,Notebook直接拼接路径
abfss://<容器名>@<存储账户名>.dfs.core.windows.net/{current_year}/Week{current_week}/Week{current_week}.xlsx进行处理 - 这种方式无需扫描所有文件,直接精准定位新增文件,效率最高,适合固定周期场景
方法四:基于事件触发的实时处理
- 如果需要准实时处理新增文件,用Azure Event Grid监听存储容器的文件创建事件
- 当有新的xlsx文件上传时,Event Grid触发Synapse管道运行
- 在管道中获取触发事件里的文件路径,传递给Notebook,Notebook直接处理该文件
- 这种方式避免定期扫描,仅在有新增数据时运行,资源利用率更高
内容的提问来源于stack exchange,提问作者HamidBee
相关产品推荐
相关产品推荐

