Spark DataFrame Schema匹配校验始终返回False问题求助
Spark Schema匹配校验返回False问题排查
开展数据治理工作时,需先校验抽取查询的Schema再执行数据转换。已实现两个类:一个用于抽取数据,另一个用于对照目标Schema校验数据。若Schema匹配则继续检查数据质量,但目前即使Schema视觉上完全匹配,代码中的if判断也始终返回False。以下是代码及运行输出:
代码实现
from pyspark.sql import SparkSession from pyspark.sql.types import * import json class Session: def __init__(self, config_file, inferred=True, header=True): self.spark = self._create_session() self.schema = inferred self.header = header self.config_file = config_file self.connections = self.parse_file() self.db_name = self.connections['db_name'] self.password = self.connections['password'] self.name = self.connections['username'] self.table_name = self.connections["tbl_name"] # self._create_session() @staticmethod def _create_session(): spark = ( SparkSession \ .builder \ .appName("test") \ .config("spark.jars", SPARK_JARS) \ .getOrCreate() ) return spark def parse_file(self): with open(self.config_file) as test: data = json.load(test) return data def pull_data(self): df = self.spark.read.format("jdbc") \ .option("url", f"jdbc:postgresql://localhost:5432/{self.db_name}") \ .option("dbtable", f'{self.table_name}') \ .option("user", f"{self.name}") \ .option("password", f"{self.password}") \ .option("driver", "org.postgresql.Driver") \ .load() return df DESIRED_SCHEMA = StructType([ StructField(name='order_id', dataType=IntegerType(), nullable=True), StructField(name='order_date', dataType=TimestampType(), nullable=True), StructField(name='order_customer_id', dataType=IntegerType(), nullable=True), StructField(name='order_status', dataType=StringType(), nullable=True) ]) class SchemaCheck: def __init__(self, targeted_schema=DESIRED_SCHEMA): self.schema = targeted_schema def val_schema(self, df): if not df.schema == self.schema: return False return True if __name__ == "__main__": valid_schema = SchemaCheck() spark = Session('connections.json') my_df = spark.pull_data() my_df.printSchema() print(my_df.schema) print(valid_schema.schema) print(valid_schema.val_schema(my_df))
运行输出
root |-- order_id: integer (nullable = true) |-- order_date: timestamp (nullable = true) |-- order_customer_id: integer (nullable = true) |-- order_status: string (nullable = true) StructType([StructField('order_id', IntegerType(), True), StructField('order_date', TimestampType(), True), StructField('order_customer_id', IntegerType(), True), StructField('order_status', StringType(), True)]) StructType([StructField('order_id', IntegerType(), True), StructField('order_date', TimestampType(), True), StructField('order_customer_id', IntegerType(), True), StructField('order_status', StringType(), True)]) False
问题原因
Spark中StructType和StructField的直接==比较,并非仅校验字段名、数据类型、可空性这些可见属性,还会检查对象的内部标识或元数据。即使两个Schema的视觉输出完全一致,它们对应的对象实例可能是独立创建的不同对象,导致相等判断返回False。
解决方法
方法1:转换为标准字符串后比较
修改SchemaCheck类的val_schema方法,利用simpleString()方法将Schema转换为统一格式的字符串再对比:
def val_schema(self, df): return df.schema.simpleString() == self.schema.simpleString()
方法2:逐个字段校验核心属性
如果需要更精细的控制(比如忽略某些非核心属性),可以遍历每个字段,分别校验名称、数据类型、可空性:
def val_schema(self, df): if len(df.schema.fields) != len(self.schema.fields): return False for df_field, target_field in zip(df.schema.fields, self.schema.fields): # 校验字段名、数据类型字符串、可空性 if (df_field.name != target_field.name or str(df_field.dataType) != str(target_field.dataType) or df_field.nullable != target_field.nullable): return False return True
内容的提问来源于stack exchange,提问作者BloodKid01
相关产品推荐
相关产品推荐

