PySpark中如何将base64图片解码后以jpg格式存入Blob容器
报错原因
Spark内置的binaryFile数据源仅实现了读取逻辑,没有提供写入能力,这是抛出Write is not supported for binary file data source的根本原因。另外内置的image数据源同样只支持读取图片元数据和二进制内容,不支持直接落地为jpg文件,所以换image格式写入也会失败。
可行方案
方案1:分布式并行写入(推荐,适配大数据量场景)
通过mapPartitions在Executor端直接调用Hadoop文件系统API写入二进制流,不需要把全量数据汇集到Driver,性能和原生Spark分布式写入一致,支持十万级以上文件批量处理。
实现代码:
from pyspark.sql.functions import col, unbase64, monotonically_increasing_id sc = spark.sparkContext hadoop_conf = sc._jsc.hadoopConfiguration() # 读取base64文本并解码,为每个文件生成唯一ID避免重名覆盖 img_rdd = spark.read.format("text").load("/mnt/pictures_text") \ .select( unbase64(col("value")).alias("content"), monotonically_increasing_id().alias("file_id") ).rdd def write_jpg_partition(rows): # 每个分区初始化一次文件系统客户端,减少重复创建开销 fs = sc._jvm.org.apache.hadoop.fs.FileSystem.get(hadoop_conf) for row in rows: save_path = sc._jvm.org.apache.hadoop.fs.Path(f"/mnt/pictures_text/copy/{row.file_id}.jpg") out = fs.create(save_path, True) out.write(row.content) out.close() return iter([]) # 触发分区写入逻辑 img_rdd.mapPartitions(write_jpg_partition).count()
使用注意:
- 提前创建目标存储路径
/mnt/pictures_text/copy,避免路径不存在报错 - 文件名规则可以自定义,如果原始数据中保留了原jpg文件名,可以替换掉示例中的自增ID
- 该方案兼容挂载到Spark的Blob Container、S3、HDFS等所有存储系统
方案2:小数据量场景:Driver端本地写入
如果总图片量在千级以内,可以把解码后的数据拉取到Driver节点直接写入,逻辑更简单:
from pyspark.sql.functions import col, unbase64 import os # 提前创建目标目录 os.makedirs("/dbfs/mnt/pictures_text/copy", exist_ok=True) # 读取、解码数据 img_data = spark.read.format("text").load("/mnt/pictures_text") \ .select(unbase64(col("value")).alias("content")) \ .collect() # 循环写入jpg for idx, row in enumerate(img_data): with open(f"/dbfs/mnt/pictures_text/copy/{idx}.jpg", "wb") as f: f.write(row.content)
注意:该方案会把所有图片二进制数据全部加载到Driver内存,数据量稍大就会触发OOM,仅适合小批量测试使用。
避坑提示
不要尝试将解码后的二进制内容转为字符串后用text/csv等文本格式写入,会因为字符编码转换导致jpg文件损坏,必须直接以字节流方式写入。
内容的提问来源于stack exchange,提问作者gorrch
相关产品推荐
相关产品推荐

