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

在(Py)Spark读取JDBC源时遭遇Unsupported Array错误的技术问询

问题:Spark批量读取PostgreSQL表时遇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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:23:42