PySpark读取Azure Blob多JSON文件过慢,如何逐个读取处理?
解决Spark读取Azure Blob大量JSON文件速度慢的问题:逐个文件处理方案
当然可以逐个读取处理啦!而且针对你的场景,这种方式还能配合一些优化点解决读取慢的问题,下面给你具体的实现思路和代码示例:
一、先搞懂为啥批量读这么慢(可选,但能帮你更精准优化)
你用spark.read.json()直接读路径慢,大概率是因为Spark默认会自动推断JSON的Schema——它得扫描所有文件的结构来确定字段类型,文件数量多、体积大时这个过程就会特别耗时。所以不管是批量还是逐个处理,提前指定Schema都是提升速度的关键。
二、逐个读取文件的两种实用方式
方式1:获取文件列表后循环读取+处理+合并
这种方式直观好调试,完全匹配你的处理需求:
步骤1:拿到Azure Blob里所有JSON文件的路径
你可以用Azure的Python SDK来列举Blob存储里的文件,代码示例如下:
from azure.storage.blob import BlobServiceClient # 替换成你的Azure Blob连接字符串和容器名 connect_str = "你的Blob连接字符串" container_name = "目标容器名" # 初始化客户端 blob_service_client = BlobServiceClient.from_connection_string(connect_str) container_client = blob_service_client.get_container_client(container_name) # 收集所有JSON文件的Spark可读路径 file_paths = [] for blob in container_client.list_blobs(): if blob.name.endswith(".json"): # 转换成wasb格式的路径,Spark能直接识别 file_path = f"wasb://{container_name}@{blob_service_client.account_name}.blob.core.windows.net/{blob.name}" file_paths.append(file_path)
如果你用的是Databricks,用dbutils.fs.ls()会更方便:
file_paths = [file.path for file in dbutils.fs.ls("wasb://你的容器路径") if file.path.endswith(".json")]
步骤2:循环处理每个文件,最后合并结果
首先一定要提前定义Schema,这能彻底避免自动推断的开销。假设你的JSON结构是包含id、content、时间字符串、数值字符串等字段,示例Schema如下:
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, TimestampType json_schema = StructType([ StructField("id", StringType(), nullable=False), StructField("content", StringType(), nullable=True), StructField("timestamp_str", StringType(), nullable=True), StructField("value_str", StringType(), nullable=True), StructField("unused_field", StringType(), nullable=True) # 你要剔除的字段 ])
然后循环处理每个文件,执行你的过滤、字段剔除、字符串拆分、类型转换逻辑:
from pyspark.sql import SparkSession from pyspark.sql.functions import split, lower, unix_timestamp, from_unixtime, col from functools import reduce from pyspark.sql import DataFrame spark = SparkSession.builder.getOrCreate() processed_dfs = [] for file_path in file_paths: # 读取单个文件,指定Schema跳过自动推断 single_file_df = spark.read.json(file_path, schema=json_schema) # 执行你的处理逻辑,按需求调整 processed_df = single_file_df \ .filter(col("id").isNotNull()) # 过滤掉id为空的记录 .drop("unused_field") # 剔除不需要的属性 .withColumn("content_lower", lower(col("content"))) # 字符串转小写 .withColumn("content_parts", split(col("content"), ",")) # 按逗号拆分字符串 .withColumn("timestamp", from_unixtime(unix_timestamp(col("timestamp_str"), "yyyy-MM-dd HH:mm:ss")).cast(TimestampType())) # 转换时间类型 .withColumn("value", col("value_str").cast(DoubleType())) # 把字符串转成数值型 .drop("timestamp_str", "value_str") # 移除原始的字符串类型字段 processed_dfs.append(processed_df) # 一次性合并所有处理后的DataFrame,比多次union更高效 final_df = reduce(DataFrame.union, processed_dfs) # 后续可以写入存储或者做其他操作 final_df.show()
小提示:
如果文件数量特别多(比如上万级),别每次循环都union,像上面那样先把所有处理后的DataFrame存到列表里,最后用reduce合并,能避免Spark执行计划变得过于复杂。
方式2:利用Spark分布式能力按文件分区处理
如果你不想手动写循环,也可以让Spark按文件路径分区,每个分区对应一个文件,然后在分区内处理数据——这种方式适合超大规模的文件集合,能利用Spark的分布式算力:
from pyspark.sql.functions import input_file_name # 读取所有文件,同时带上文件路径字段,指定Schema df_with_path = spark.read.json("wasb://你的容器路径", schema=json_schema) \ .withColumn("file_path", input_file_name()) # 按文件路径分区,确保每个分区对应单个文件 df_with_path = df_with_path.repartition(col("file_path")) # 定义分区处理函数,这里可以实现和之前一样的逻辑 def process_partition(partition): for row in partition: # 这里可以对每行数据做处理,比如过滤、字段转换等 # 处理后返回新的行对象 processed_row = ( row.id, row.content.lower(), row.content.split(","), from_unixtime(unix_timestamp(row.timestamp_str, "yyyy-MM-dd HH:mm:ss")).cast(TimestampType()), float(row.value_str) if row.value_str else None ) yield processed_row # 应用分区处理,转成DataFrame processed_df = df_with_path.rdd.mapPartitions(process_partition).toDF( ["id", "content_lower", "content_parts", "timestamp", "value"] )
这种方式调试起来不如第一种直观,但分布式处理效率更高,适合超大规模场景。
三、额外优化小技巧
- 必须指定Schema:再强调一次,这是提升读取速度最有效的手段,没有之一。
- 合并小文件:如果你的JSON都是几十KB的小文件,建议先在Blob存储里合并成100MB-1GB左右的大文件,Spark处理大文件的效率远高于大量小文件。
- 换用ABFSS协议:如果你的Azure存储支持Data Lake Storage Gen2,尽量用
abfss://代替wasb://,ABFSS是更高效的协议,性能提升明显。
内容的提问来源于stack exchange,提问作者Jiew Meng
相关产品推荐
相关产品推荐

