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

PyFlink查询含jsonb列PostgreSQL表的JDBC报错及类型转换问题求助

问题场景

1. 初始查询报错

使用PyFlink 1.16搭配flink-connector-jdbc-1.16.0.jar查询包含jsonb列的PostgreSQL表,即使SQL未涉及jsonb列,执行查询时仍触发类型不支持错误:

sql = "SELECT entity_id FROM event_files"
table2 = table_env.sql_query(sql)
table2.execute().print()

报错信息:

Caused by: java.lang.UnsupportedOperationException: Doesn't support Postgres type 'jsonb' yet
    at org.apache.flink.connector.jdbc.dialect.psql.PostgresTypeMapper.mapping(PostgresTypeMapper.java:173)

2. 修改类型映射后新报错

将PostgresTypeMapper.java中添加jsonb到String的映射:

private static final String PG_CHARACTER_VARYING_ARRAY = "_varchar";
private static final String PG_JSONB = "jsonb";

// 在mapping方法的case分支中添加
case PG_DATE_ARRAY:
    return DataTypes.ARRAY(DataTypes.DATE());
case PG_JSONB:
    return DataTypes.STRING();

重新编译后,基础查询可正常运行,但当查询并转换jsonb列时,出现类型转换异常:

sql = "SELECT entity_id, CAST(payload AS varchar) AS P1 FROM event_files"
table2 = table_env.sql_query(sql)
table2.execute().print()

报错信息:

Caused by: java.lang.ClassCastException: class org.postgresql.util.PGobject cannot be cast to class java.lang.String (org.postgresql.util.PGobject is in unnamed module of loader org.apache.flink.util.ChildFirstClassLoader @38c50913; java.lang.String is in module java.base of loader 'bootstrap')
    at org.apache.flink.connector.jdbc.converter.AbstractJdbcRowConverter.lambda$createInternalConverter$224afae6$10(AbstractJdbcRowConverter.java:176)

解决方案

要彻底解决问题,需要同时处理类型映射和数据转换逻辑:

步骤1:完善PostgresTypeMapper的类型映射(已完成)

确保PostgresTypeMapper中正确将jsonb类型映射为Flink的STRING类型,即之前添加的代码:

private static final String PG_JSONB = "jsonb";

case PG_JSONB:
    return DataTypes.STRING();

步骤2:修改PostgresRowConverter,处理PGobject到String的转换

Flink读取PostgreSQL的jsonb数据时,驱动返回的是PGobject对象,默认转换器无法直接将其转为String,需要在PostgresRowConverter中添加专属转换逻辑:

找到PostgresRowConverter的createInternalConverter方法,补充针对STRING类型的转换器:

// 在createInternalConverter方法中添加
if (logicalType instanceof VarCharType) {
    return (val) -> {
        // 处理jsonb对应的PGobject对象
        if (val instanceof PGobject) {
            return ((PGobject) val).getValue();
        }
        return val;
    };
}

步骤3:重新编译替换jar包

重新编译flink-connector-jdbc模块,将生成的新jar包替换原有的flink-connector-jdbc-1.16.0.jar。

备选方案:SQL层面转换(无需修改源码)

如果不想修改Flink源码,可以直接在SQL中使用PostgreSQL内置的类型转换语法,将jsonb转为文本:

sql = "SELECT entity_id, payload::text AS P1 FROM event_files"
# 或者使用jsonb_to_text函数(部分PostgreSQL版本支持)
# sql = "SELECT entity_id, jsonb_to_text(payload) AS P1 FROM event_files"

原因说明

  • 初始报错是因为Flink JDBC连接器在解析表结构时,会遍历所有列的类型,即使查询未涉及jsonb列,仍会触发类型映射检查,而默认未支持jsonb类型。
  • 修改TypeMapper后,类型映射已生效,但数据读取时,PostgreSQL驱动返回的PGobject无法被默认转换器直接转为String,因此需要显式提取PGobject的字符串值。

内容的提问来源于stack exchange,提问作者Sami Badawi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 22:31:59