使用NiFi PutDatabaseRecord向PostgreSQL jsonb字段插数据失败排查
问题描述
环境信息:
- Apache NiFi版本:nifi-2.0.0-M2-RC4
- PostgreSQL版本:12.12
尝试使用AVROSchemaRegistry和PutDatabaseRecord处理器,将包含JSON结构的记录插入PostgreSQL的jsonb类型字段时失败。
AVRO Schema
{ "name": "ebs_data", "namespace": "nifi", "type": "record", "fields": [ { "***": "***" }, { "name": "log", "type": { "type": "array", "items": { "type": "record", "name": "log_rec", "fields": [ {"name": "log_level", "type": ["string","null"]}, {"name": "log_message", "type": ["string","null"]}, {"name": "log_date", "type": ["string","null"]} ] } } } ] }
错误堆栈
2024-07-29 12:50:29,701 ERROR [Timer-Driven Process Thread-2] o.a.n.p.standard.PutDatabaseRecord PutDatabaseRecord[id=1f964dda-cd85-3499-190a-a3724f133e23] Failed to put Records to database for StandardFlowFileRecord[uuid=acbcce68-f9f4-45b3-b88d-7e8694167097,claim=StandardContentClaim [resourceClaim=StandardResourceClaim[id=1722 246429765-23, container=default, section=23], offset=9744, length=73727],offset=0,name=349045823834552,size=73727]. Routing to failure. java.sql.BatchUpdateException: Batch entry 0 INSERT INTO ebs_data.ebs_data.rnd_heap_and_parts (***, log) VALUES (***,('[Ljava.lang.Object;@1d1458e1')) was aborted: ERROR: invalid input syntax for type json Detail: Token "Ljava" is invalid. Where: JSON data, line 1: [Ljava... Call getNextException to see other errors in the batch. at org.postgresql.jdbc.BatchResultHandler.handleError(BatchResultHandler.java:165) at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2413) at org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:579) at org.postgresql.jdbc.PgStatement.internalExecuteBatch(PgStatement.java:912) at org.postgresql.jdbc.PgStatement.executeBatch(PgStatement.java:936) at org.postgresql.jdbc.PgPreparedStatement.executeBatch(PgPreparedStatement.java:1733) at com.zaxxer.hikari.pool.ProxyStatement.executeBatch(ProxyStatement.java:127) at com.zaxxer.hikari.pool.HikariProxyPreparedStatement.executeBatch(HikariProxyPreparedStatement.java) at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(DirectMethodHandleAccessor.java:103) at java.base/java.lang.reflect.Method.invoke(Method.java:580) at org.apache.nifi.controller.service.StandardControllerServiceInvocationHandler.invoke(StandardControllerServiceInvocationHandler.java:254) at org.apache.nifi.controller.service.StandardControllerServiceInvocationHandler$ProxiedReturnObjectInvocationHandler.invoke(StandardControllerServiceInvocationHandler.java:240) at jdk.proxy23/jdk.proxy23.$Proxy206.executeBatch(Unknown Source) at org.apache.nifi.processors.standard.PutDatabaseRecord.executeDML(PutDatabaseRecord.java:958) at org.apache.nifi.processors.standard.PutDatabaseRecord.putToDatabase(PutDatabaseRecord.java:1140) at org.apache.nifi.processors.standard.PutDatabaseRecord.onTrigger(PutDatabaseRecord.java:559) at org.apache.nifi.processor.AbstractProcessor.onTrigger(AbstractProcessor.java:27) at org.apache.nifi.controller.StandardProcessorNode.onTrigger(StandardProcessorNode.java:1274) at org.apache.nifi.controller.tasks.ConnectableTask.invoke(ConnectableTask.java:244) at org.apache.nifi.controller.scheduling.TimerDrivenSchedulingAgent$1.run(TimerDrivenSchedulingAgent.java:102) at org.apache.nifi.engine.FlowEngine$2.run(FlowEngine.java:110) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:572) at java.base/java.util.concurrent.FutureTask.runAndReset(FutureTask.java:358) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:305) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) at java.base/java.lang.Thread.run(Thread.java:1583) Caused by: org.postgresql.util.PSQLException: ERROR: invalid input syntax for type json Detail: Token "Ljava" is invalid. Where: JSON data, line 1: [Ljava... at org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2725) at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2412) ... 25 common frames omitted
问题原因与解决方案
原因
错误日志中的[Ljava.lang.Object;@1d1458e1是Java数组对象的默认toString输出,说明PutDatabaseRecord处理器未将AVRO的复杂数组类型正确序列化为JSON字符串,而是直接传入了Java对象的字符串形式,导致PostgreSQL的jsonb字段语法校验失败。
解决步骤
配置字段类型映射
- 打开PutDatabaseRecord处理器配置,在
Database Type Name Mapping属性中添加:
该配置会告知处理器将log=jsonblog字段按PostgreSQL的jsonb类型处理,自动完成JSON序列化。
- 打开PutDatabaseRecord处理器配置,在
调整Record Reader配置
- 检查关联的AVRO Record Reader控制器服务,确保开启了复杂类型的JSON序列化支持。如果使用
AvroRecordReader,确认其配置中没有禁用类型转换;若原始数据是JSON格式,可改用JsonRecordReader直接读取JSON数据。
- 检查关联的AVRO Record Reader控制器服务,确保开启了复杂类型的JSON序列化支持。如果使用
提前转换为JSON字符串(备选)
- 使用
ConvertRecord处理器,搭配AvroRecordReader和JsonRecordSetWriter,将AVRO数据转换为JSON格式,此时log字段会被序列化为合法的JSON数组字符串。再将转换后的JSON数据传入PutDatabaseRecord插入到jsonb字段。
- 使用
SQL显式类型转换(备选)
- 若上述方法无效,可在PutDatabaseRecord的
SQL Statement中强制转换类型:
需确保传入的INSERT INTO ebs_data.ebs_data.rnd_heap_and_parts (***, log) VALUES (***, CAST(? AS jsonb))log参数已经是合法的JSON字符串。
- 若上述方法无效,可在PutDatabaseRecord的
检查JDBC驱动版本
- 确认NiFi使用的PostgreSQL JDBC驱动版本与12.12兼容,建议使用42.2.x及以上版本,避免类型转换兼容性问题。
内容的提问来源于stack exchange,提问作者Будён Михайлович Семённый
相关产品推荐
相关产品推荐

