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

如何在Spark DataFrame中新增/更新列并实现数据校验?

解决方案

一、生成目标校验结果DataFrame

核心思路是遍历每个表的路径,读取数据后执行各类校验,将每个表的校验结果整理成单行数据,最终合并成目标结构的DataFrame。以下是具体实现:

1. 导入依赖

from pyspark.sql import functions as F, types as T

2. 定义单表处理函数

编写函数完成单个表的校验逻辑,返回包含所有校验结果的元组:

def process_single_table(s3_prefix, table_name, bucket_name):
    # 读取对应路径的Parquet文件
    s3_path = f"s3://{bucket_name}/{s3_prefix}/*.parquet"
    ing_df = spark.read.parquet(s3_path)
    
    # 1. 获取ID长度(假设该表ID长度统一,若需校验一致性可扩展逻辑)
    id_length = ing_df.select(F.length('ID')).first()[0]
    
    # 2. 检查description是否存在换行符等特殊字符
    has_strange_char = ing_df.filter(F.col('description').rlike(r'[\n\r]')).count() > 0
    strange_char_flag = "YES" if has_strange_char else "NO"
    
    # 3. 收集异常ID(示例:长度不符合标准或包含非字母数字字符)
    bad_ids_df = ing_df.filter(
        F.length('ID') != id_length 
        | F.col('ID').rlike(r'[^a-zA-Z0-9]')
    ).select(F.collect_set('ID').alias('bad_ids')).first()
    bad_ids_list = bad_ids_df['bad_ids'] if bad_ids_df['bad_ids'] else []
    
    # 返回当前表的校验结果
    return (s3_prefix, table_name, id_length, strange_char_flag, bad_ids_list)

3. 批量处理生成最终DataFrame

遍历初始的df_result,调用处理函数收集所有结果,再转换为目标结构的DataFrame:

# 收集所有表的校验结果
validation_results = []
for row in df_result.collect():
    res = process_single_table(row['s3_prefix'], row['table_name'], bucket_name)
    validation_results.append(res)

# 定义目标Schema
final_schema = T.StructType([
    T.StructField("s3_prefix", T.StringType(), nullable=True),
    T.StructField("table_name", T.StringType(), nullable=True),
    T.StructField("id_length", T.IntegerType(), nullable=True),
    T.StructField("strange_character", T.StringType(), nullable=True),
    T.StructField("id_strange_character", T.ArrayType(T.StringType()), nullable=True)
])

# 创建最终校验结果DataFrame
final_validation_df = spark.createDataFrame(validation_results, schema=final_schema)
final_validation_df.show(truncate=False)

二、正则检查description字段的换行符

要检测description中是否存在换行符(包括\n、\r、\r\n),使用Spark的rlike函数配合正则表达式[\n\r]即可,该表达式可覆盖绝大多数换行场景:

# 检查单条记录是否包含换行符
has_newline = F.col('description').rlike(r'[\n\r]')

# 统计表中存在换行符的记录总数
newline_record_count = ing_df.filter(has_newline).count()

如果需要精确匹配Windows风格的\r\n换行符,可将正则替换为r'\r\n'。


内容的提问来源于stack exchange,提问作者Julio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 04:02:25