在(Py)Spark读取JDBC源时遭遇Unsupported Array错误的技术问询
我尝试用Spark把PostgreSQL数据库里的所有表转成DataFrame,写了如下代码:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("Connect to DB") \ .getOrCreate() jdbcUrl = "jdbc:postgresql://XXXXXX" connectionProperties = { "user" : " ", "password" : " ", "driver" : "org.postgresql.Driver" } query = "(SELECT table_name FROM information_schema.tables) XXX" df = spark.read.jdbc(url=jdbcUrl, table=query, properties=connectionProperties) table_name_list = df.select("table_name").rdd.flatMap(lambda x: x).collect() for table_name in table_name_list: df2 = spark.read.jdbc(url=jdbcUrl, table=table_name, properties=connectionProperties)
结果运行时抛出错误:
java.sql.SQLException: Unsupported type ARRAY on generating df2 for table name
但如果我直接硬编码表名(比如下面这样)就完全没问题:
df2 = spark.read.jdbc(jdbcUrl,"conditions",properties=connectionProperties)
我已经确认table_name的类型是String,想问问当前的批量读取方法是否正确?
你的批量读取思路本身是对的——先捞全库表名再循环读取,但踩坑的点在于:你循环读取的表中,有部分包含PostgreSQL的ARRAY类型字段,而Spark JDBC默认没有做好这种类型的映射处理。你硬编码时没报错,只是刚好选的conditions表没有ARRAY类型字段而已。
给你几个实用的解决方向:
1. 让Spark自动将ARRAY转成字符串处理
在connectionProperties里添加stringtype=unspecified参数,这样Spark会把它识别不了的类型(比如ARRAY)统一转成字符串,后续你可以在Spark里再根据业务需求解析这些字符串:
connectionProperties = { "user" : "你的用户名", "password" : "你的密码", "driver" : "org.postgresql.Driver", "stringtype": "unspecified" # 新增这个参数 }
2. 过滤掉含ARRAY类型的表
如果这些带ARRAY类型的表你暂时不需要,可以先在查询表名时就过滤掉它们,只读取不含ARRAY字段的表:
# 修改查询语句,只获取不含ARRAY类型字段的表 query = """ (SELECT DISTINCT table_name FROM information_schema.columns WHERE data_type != 'ARRAY') as filtered_tables """ df = spark.read.jdbc(url=jdbcUrl, table=query, properties=connectionProperties) table_name_list = df.select("table_name").rdd.map(lambda x: x[0]).collect()
3. 手动指定Schema保留ARRAY类型
如果需要保留ARRAY类型的数据,可以提前为每个表定义对应的Spark Schema,读取时通过schema参数指定,把PostgreSQL的ARRAY映射成Spark的ArrayType。不过这种方式需要逐个处理表结构,更适合表数量不多的场景。
另外给你个小优化:你用来提取表名的df.select("table_name").rdd.flatMap(lambda x: x).collect()可以简化成df.select("table_name").rdd.map(lambda x: x[0]).collect(),如果表名数量不多的话,用df.toPandas()['table_name'].tolist()会更高效哦。
内容的提问来源于stack exchange,提问作者K.Pil

