使用Spark JDBC读取Redshift时为何报Parquet schema推断错误?
问题原因分析
这个错误其实是Spark 2.x的一个典型“默认行为坑”:当你调用spark.read.load()却没明确指定数据源格式时,Spark会默认把输入当成Parquet文件来处理——哪怕你配置了JDBC的连接参数,Spark依然会忽略这些参数,按默认的Parquet逻辑去尝试推断schema,这就导致了明明是连接Redshift却抛出Parquet相关错误的矛盾情况。
两种解决办法
你只需要明确告诉Spark要使用JDBC数据源即可,有两种常用实现方式:
方式1:添加.format("jdbc")指定数据源
在spark.read之后加上.format("jdbc"),让Spark知道要走JDBC连接逻辑,修改后的代码如下:
val query = "SELECT * FROM some_big_table WHERE something > 1" val df : DataFrame = spark.read .format("jdbc") // 关键:指定使用JDBC数据源 .option("url", s"""jdbc:postgresql://${redshiftInfo.hostnameAndPort}/${redshiftInfo.database}?currentSchema=${redshiftInfo.schema}""" ) .option("user", redshiftInfo.username) .option("password", redshiftInfo.password) .option("dbtable", query) .load()
方式2:直接使用spark.read.jdbc()专用方法
这种方式更直观,不需要额外指定format,代码结构也更清晰,注意查询语句需要包装成子查询并命名临时表(JDBC数据源的要求):
// 把查询包装成子查询,命名为临时表 val query = "(SELECT * FROM some_big_table WHERE something > 1) AS temp_table" val jdbcUrl = s"""jdbc:postgresql://${redshiftInfo.hostnameAndPort}/${redshiftInfo.database}?currentSchema=${redshiftInfo.schema}""" val connectionProps = new java.util.Properties() connectionProps.setProperty("user", redshiftInfo.username) connectionProps.setProperty("password", redshiftInfo.password) val df: DataFrame = spark.read.jdbc(jdbcUrl, query, connectionProps)
补充说明
Spark 2.x把Parquet设为默认数据源是为了优化大数据场景下的读写性能,但在连接关系型数据库这类非文件型数据源时,必须显式指定jdbc格式,否则就会触发默认的Parquet解析逻辑,出现你遇到的这类“驴唇不对马嘴”的错误。
内容的提问来源于stack exchange,提问作者hotmeatballsoup
相关产品推荐
相关产品推荐

