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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 04:03:23