PyFlink查询含jsonb列PostgreSQL表的JDBC报错及类型转换问题求助
解决PyFlink 1.16查询PostgreSQL jsonb列的类型转换问题
问题场景
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
相关产品推荐
相关产品推荐

