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

如何在Databricks的PySpark流作业中创建ADLS Gen2目录

问题重现(模拟代码)
from pyspark.sql import SparkSession
import cv2
import numpy as np
from pyspark.sql.functions import udf, col

spark = SparkSession.builder.appName("ImageStreamWrite").getOrCreate()

# 模拟含图片二进制、图像属性的流数据
stream_df = spark.readStream.format("kafka")\
    .option("kafka.bootstrap.servers", "xxx:9092")\
    .load()

def write_image(img_bytes, category, resolution):
    # 按属性构建ADLS Gen2目标路径
    target_dir = f"abfss://container@storageaccount.dfs.core.windows.net/images/{category}/{resolution}"
    target_path = f"{target_dir}/img.jpg"
    
    # 尝试用dbutils创建目录会报错:dbutils未定义(Executor端无法调用Driver端工具)
    # dbutils.fs.mkdirs(target_dir)
    
    # 目录不存在时cv2写入失败
    img = cv2.imdecode(np.frombuffer(img_bytes, np.uint8), cv2.IMREAD_COLOR)
    cv2.imwrite(target_path, img)
    return "success"

write_image_udf = udf(write_image)

query = stream_df.withColumn("status", write_image_udf(col("value"), col("category"), col("resolution")))\
    .writeStream.format("console").start()

query.awaitTermination()
解决方案

因为PySpark流作业的UDF运行在Executor节点,而dbutils是Driver端专属工具,无法跨节点调用。以下两种方案可解决目录创建问题:

方案1:Executor端用Hadoop FileSystem API动态创建目录

通过Hadoop原生API直接操作ADLS Gen2路径,支持在Executor节点检查并创建目录:

from pyspark.sql import SparkSession
import cv2
import numpy as np
from pyspark.sql.functions import udf, col
from py4j.java_gateway import java_import
from pyspark.context import SparkContext

def get_hadoop_fs(sc):
    # 导入Hadoop路径和文件系统类
    java_import(sc._jvm, "org.apache.hadoop.fs.Path")
    java_import(sc._jvm, "org.apache.hadoop.fs.FileSystem")
    return sc._jvm.FileSystem.get(sc._jsc.hadoopConfiguration())

def write_image(img_bytes, category, resolution):
    target_dir = f"abfss://container@storageaccount.dfs.core.windows.net/images/{category}/{resolution}"
    target_path = f"{target_dir}/img.jpg"
    
    # 获取Executor端的Hadoop文件系统实例
    sc = SparkContext.getOrCreate()
    fs = get_hadoop_fs(sc)
    hadoop_dir_path = sc._jvm.Path(target_dir)
    
    # 检查目录是否存在,不存在则递归创建
    if not fs.exists(hadoop_dir_path):
        fs.mkdirs(hadoop_dir_path)
    
    # 写入图片
    img = cv2.imdecode(np.frombuffer(img_bytes, np.uint8), cv2.IMREAD_COLOR)
    cv2.imwrite(target_path, img)
    return "success"

write_image_udf = udf(write_image)

注意事项

  • 确保所有Executor节点安装了opencv-python,可通过集群初始化脚本执行:pip install opencv-python
  • 作业所用身份(服务主体/MSI)需拥有ADLS Gen2的Storage Blob Data Contributor权限
  • 必须使用abfss://<容器名>@<存储账户>.dfs.core.windows.net/格式的原生路径,避免DBFS挂载路径的解析问题

方案2:Driver端提前批量创建目录(适合属性固定场景)

如果图像的category、resolution是有限的已知值,可在Driver端提前用dbutils批量创建所有可能的目录,避免Executor端重复操作:

# 在Driver端执行(流作业启动前)
categories = ["cat", "dog", "bird"]
resolutions = ["1080p", "720p", "480p"]

for cat in categories:
    for res in resolutions:
        dir_path = f"abfss://container@storageaccount.dfs.core.windows.net/images/{cat}/{res}"
        dbutils.fs.mkdirs(dir_path)

内容的提问来源于stack exchange,提问作者Divyam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 16:05:28