如何在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
相关产品推荐
相关产品推荐

