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

使用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字段语法校验失败。

解决步骤

  1. 配置字段类型映射

    • 打开PutDatabaseRecord处理器配置,在Database Type Name Mapping属性中添加:
      log=jsonb
      
      该配置会告知处理器将log字段按PostgreSQL的jsonb类型处理,自动完成JSON序列化。
  2. 调整Record Reader配置

    • 检查关联的AVRO Record Reader控制器服务,确保开启了复杂类型的JSON序列化支持。如果使用AvroRecordReader,确认其配置中没有禁用类型转换;若原始数据是JSON格式,可改用JsonRecordReader直接读取JSON数据。
  3. 提前转换为JSON字符串(备选)

    • 使用ConvertRecord处理器,搭配AvroRecordReader和JsonRecordSetWriter,将AVRO数据转换为JSON格式,此时log字段会被序列化为合法的JSON数组字符串。再将转换后的JSON数据传入PutDatabaseRecord插入到jsonb字段。
  4. SQL显式类型转换(备选)

    • 若上述方法无效,可在PutDatabaseRecord的SQL Statement中强制转换类型:
      INSERT INTO ebs_data.ebs_data.rnd_heap_and_parts (***, log) VALUES (***, CAST(? AS jsonb))
      
      需确保传入的log参数已经是合法的JSON字符串。
  5. 检查JDBC驱动版本

    • 确认NiFi使用的PostgreSQL JDBC驱动版本与12.12兼容,建议使用42.2.x及以上版本,避免类型转换兼容性问题。

内容的提问来源于stack exchange,提问作者Будён Михайлович Семённый

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 08:57:19