Databricks Autoloader写入含非法字符嵌套列的问题解决
修复Databricks Autoloader写入时嵌套列的非法字符问题
使用Databricks Autoloader写入数据时,若嵌套列名称包含非法字符,会触发如下报错:
Found invalid character(s) among " ,;{}()\n\t=" in the column names of your schema.
已知仅能通过如下代码修复外层列,但无法处理未展开的嵌套列:
for col in df.columns: df = df.select([col(c).alias(re.sub("[^0-9a-zA-Z\_]+","",c)) for c in df.columns])
解决方案:递归处理所有层级的列名
要处理嵌套列,需要递归遍历Schema的所有字段,替换每个字段名称中的非法字符,同时保留原有的嵌套结构。以下是实现代码:
首先导入必要的库:
import re from pyspark.sql.types import StructType, StructField, ArrayType, MapType
定义递归清理Schema的函数:
def clean_schema(schema): def clean_field_name(name): # 替换非法字符为空,仅保留字母、数字和下划线 return re.sub("[^0-9a-zA-Z\_]+", "", name) cleaned_fields = [] for field in schema.fields: cleaned_name = clean_field_name(field.name) # 处理Struct类型的嵌套字段 if isinstance(field.dataType, StructType): cleaned_data_type = clean_schema(field.dataType) cleaned_fields.append(StructField(cleaned_name, cleaned_data_type, field.nullable)) # 处理Array类型,若元素是Struct则递归处理 elif isinstance(field.dataType, ArrayType): if isinstance(field.dataType.elementType, StructType): cleaned_element_type = clean_schema(field.dataType.elementType) cleaned_data_type = ArrayType(cleaned_element_type, field.dataType.containsNull) else: cleaned_data_type = field.dataType cleaned_fields.append(StructField(cleaned_name, cleaned_data_type, field.nullable)) # 处理Map类型,仅递归处理值为Struct的情况 elif isinstance(field.dataType, MapType): if isinstance(field.dataType.valueType, StructType): cleaned_value_type = clean_schema(field.dataType.valueType) cleaned_data_type = MapType(field.dataType.keyType, cleaned_value_type, field.dataType.valueContainsNull) else: cleaned_data_type = field.dataType cleaned_fields.append(StructField(cleaned_name, cleaned_data_type, field.nullable)) # 基础类型直接保留 else: cleaned_fields.append(StructField(cleaned_name, field.dataType, field.nullable)) return StructType(cleaned_fields)
定义递归重命名DataFrame所有列的函数:
def rename_nested_columns(df): def rename_col(col_name): return re.sub("[^0-9a-zA-Z\_]+", "", col_name) def process_column(col): col_name = col._jc.toString().split(".")[-1] cleaned_name = rename_col(col_name) # 处理Struct类型列 if col.dataType.typeName() == "struct": renamed_fields = [process_column(col[field]).alias(rename_col(field)) for field in col.dataType.names] return col.alias(cleaned_name).struct(*renamed_fields) # 处理Array类型列,若元素是Struct则递归处理 elif col.dataType.typeName() == "array": element_type = col.dataType.elementType if element_type.typeName() == "struct": processed_element = process_column(col.getItem(0)) return col.alias(cleaned_name).array(processed_element) else: return col.alias(cleaned_name) # 处理Map类型列,仅重命名值为Struct的情况 elif col.dataType.typeName() == "map": value_type = col.dataType.valueType if value_type.typeName() == "struct": processed_value = process_column(col.getItem("key")) return col.alias(cleaned_name).map(col.keys(), processed_value) else: return col.alias(cleaned_name) # 基础类型直接重命名 else: return col.alias(cleaned_name) # 处理所有顶层列 processed_columns = [process_column(df[col]) for col in df.columns] return df.select(*processed_columns)
最后调用函数处理DataFrame:
# 清理Schema(确保与重命名后的列结构一致) cleaned_schema = clean_schema(df.schema) # 重命名所有层级的列 cleaned_df = rename_nested_columns(df) # 现在可以用cleaned_df进行Autoloader写入操作
说明
- 该方案会递归遍历所有层级的嵌套结构(Struct、Array中的Struct、Map中的Struct),将所有列名中的非法字符替换为空,仅保留字母、数字和下划线。
- 处理后的DataFrame Schema与列名完全符合Databricks的要求,不会再触发非法字符报错。
内容的提问来源于stack exchange,提问作者Preben Brudvik Olsen
相关产品推荐
相关产品推荐

