PySpark清理字符串列失败:DB2数据写入Parquet遇字符转换错误
解决DB2转Spark DataFrame时的CharConversionException问题
你的问题核心是错误发生在数据读取阶段,而非后续的Spark清理阶段。当前代码先通过JDBC加载整张表,但DB2的CHAR(8)类型USER_ID列中包含无法转成UTF-8的非法字符,JDBC驱动在读取时直接抛出CharConversionException,根本没机会执行你写的regexp_replace清理逻辑。
方案1:修改JDBC驱动参数,自动处理非法字符
DB2 JDBC驱动提供了参数控制字符转换的容错性,你可以在JDBC URL中添加db2.jcc.charsetDecoderEncoder=3,让驱动自动将非法字符替换为问号,这样就能正常读取数据,之后再在Spark中清理:
# 修改DB2 URL,添加字符处理参数 db2_url = f"{db2_url};db2.jcc.charsetDecoderEncoder=3" df = spark.read.format("jdbc") \ .option("url", db2_url) \ .option("driver", "com.ibm.db2.jcc.DB2Driver") \ .option("dbtable", db_table) \ .option("user", db2_user) \ .option("password", db2_password) \ .option("charset", "UTF-8") \ .load() # 后续清理逻辑不变 df_cleaned = df.withColumn("USER_ID", regexp_replace("USER_ID", "[^a-zA-Z0-9]", "")) df_cleaned.write.mode("append").option("maxRecordsPerFile", 5000000).parquet(hdfs_path)
参数说明:db2.jcc.charsetDecoderEncoder的取值:
- 1:默认,使用JVM默认的字符解码器
- 2:严格模式,遇到非法字符直接抛出异常(就是你当前的情况)
- 3:替换模式,将非法字符替换为
?
方案2:动态生成JDBC查询,在DB2端预处理数据
如果不想依赖驱动参数,可以在读取时动态构造带清理逻辑的查询,自动识别包含问题列的表,避免硬编码:
步骤1:自动检测表是否包含CHAR(8)类型的USER_ID列
通过DB2系统表SYSCAT.COLUMNS查询表结构,判断是否存在目标列:
import jaydebeapi def has_problematic_user_id(db_table, db2_url, db2_user, db2_password): # 建立JDBC连接查询元数据 conn = jaydebeapi.connect( "com.ibm.db2.jcc.DB2Driver", db2_url, [db2_user, db2_password] ) cursor = conn.cursor() # DB2表名、列名默认是大写,转成大写匹配 cursor.execute(f""" SELECT 1 FROM SYSCAT.COLUMNS WHERE TABNAME = '{db_table.upper()}' AND COLNAME = 'USER_ID' AND TYPENAME = 'CHAR' AND LENGTH = 8 """) result = cursor.fetchone() cursor.close() conn.close() return result is not None
步骤2:动态构造读取逻辑
根据检测结果,决定直接读表还是读带清理的子查询:
connection_props = { "user": db2_user, "password": db2_password, "driver": "com.ibm.db2.jcc.DB2Driver" } if has_problematic_user_id(db_table, db2_url, db2_user, db2_password): # 构造带清理的子查询,替换USER_ID后保留其他列 query = f""" (SELECT REGEXP_REPLACE(USER_ID, '[^\\\\x20-\\\\x7E]', '') AS USER_ID, t.* EXCLUDE(USER_ID) FROM {db_table} t) AS cleaned_table """ df = spark.read.jdbc(url=db2_url, table=query, properties=connection_props) else: # 正常读取整张表 df = spark.read.jdbc(url=db2_url, table=db_table, properties=connection_props) # 后续写入逻辑不变 df.write.mode("append").option("maxRecordsPerFile", 5000000).parquet(hdfs_path)
这种方式的好处是在DB2端完成清理,避免驱动读取时的字符转换错误,同时自动适配所有表,不用手动处理20多张表。
为什么你原来的方法无效?
你的代码是先读取整张表,再执行清理,但错误发生在JDBC驱动读取数据并转换字符的阶段——驱动在把DB2的CHAR列数据转换成Spark的String类型时,遇到了无法识别的非法字符,直接抛出异常,根本没进入到Spark的regexp_replace步骤。
内容的提问来源于stack exchange,提问作者Aaronel10
相关产品推荐
相关产品推荐

