如何验证嵌套JSON文件与DataFrame的记录数一致性?
嵌套JSON全量加载至DataFrame的验证方案
一、Python统计原始嵌套JSON的记录数
可以直接读取Blob存储中的JSON文件,通过递归遍历结构统计符合业务定义的记录数。以你提供的示例JSON为例,假设数组中的每个元素、独立的对象(如country)都算作一条记录,代码实现如下:
import json from azure.storage.blob import BlobServiceClient # 配置Blob存储连接信息 connection_string = "你的Blob连接字符串" container_name = "目标容器名" blob_name = "嵌套JSON文件的路径" # 读取Blob中的JSON文件 blob_service_client = BlobServiceClient.from_connection_string(connection_string) blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name) json_content = blob_client.download_blob().readall() json_data = json.loads(json_content) # 递归统计记录数 def count_target_records(data): count = 0 if isinstance(data, list): # 数组每个元素算一条记录 count += len(data) for item in data: count += count_target_records(item) elif isinstance(data, dict): for value in data.values(): # 无嵌套的独立对象算一条记录 if isinstance(value, dict) and not any(isinstance(v, (list, dict)) for v in value.values()): count += 1 count += count_target_records(value) return count # 计算预期总记录数 expected_total = count_target_records(json_data) print(f"原始JSON预期记录数: {expected_total}")
针对你的示例JSON,这段代码会统计出6条记录(coffee.region2条 + coffee.country1条 + brewing.region2条 + brewing.country1条)。
二、Databricks中验证DataFrame记录数
在Databricks中,需要先将嵌套JSON的结构扁平化,再统计对应记录数,和预期值对比:
1. 读取JSON文件到DataFrame
# 读取Blob存储中的JSON(需提前配置存储访问权限) df = spark.read.json("abfss://<容器名>@<存储账户名>.dfs.core.windows.net/<JSON文件路径>")
2. 拆分嵌套结构并统计记录数
根据业务定义的“记录”类型,拆分统计:
from pyspark.sql.functions import col, explode # 统计coffee.region的记录数 coffee_region_count = df.select(explode(col("coffee.region"))).count() # 统计coffee.country的记录数(独立对象算1条) coffee_country_count = df.select(col("coffee.country")).distinct().count() # 同理统计brewing下的记录数 brewing_region_count = df.select(explode(col("brewing.region"))).count() brewing_country_count = df.select(col("brewing.country")).distinct().count() # 计算实际加载的总记录数 actual_total = coffee_region_count + coffee_country_count + brewing_region_count + brewing_country_count
3. 断言验证
通过断言判断实际加载数和预期数是否一致:
assert actual_total == expected_total, f"记录加载不完整:预期{expected_total}条,实际{actual_total}条"
三、额外验证手段
除了记录数对比,还可以通过以下方式确保数据完整性:
- 校验Schema一致性:打印DataFrame的Schema,确认字段结构和原始JSON匹配
df.printSchema() - 抽样内容验证:抽取部分数据对比原始JSON和DataFrame中的内容
# 验证coffee.region第一条记录的id和name df_coffee_region_first = df.select(col("coffee.region")[0].id, col("coffee.region")[0].name).first() assert df_coffee_region_first.id == json_data["coffee"]["region"][0]["id"] assert df_coffee_region_first.name == json_data["coffee"]["region"][0]["name"]
内容的提问来源于stack exchange,提问作者Swati B
相关产品推荐
相关产品推荐

