使用Kafka Sink Connector写入Oracle遇ORA-01400空值错误求助
我是Kafka Sink Connector新手,用它往Oracle数据库写入数据时触发ORA-01400错误,提示无法向TEST1.RENTED_PRODUCT.RENTED_PRODUCT_NO插入空值。该字段是Oracle表的主键,原本通过JPA配置序列自动生成,但Kafka消息的schema和payload中均未包含这个字段。在MySQL环境下运行正常,切换至Oracle后出现此问题,恳请协助排查解决。
错误日志
Error : 1400, Position : 0, SQL = INSERT INTO "RENTED_PRODUCT"("PRODUCT_NO","RENT_START_DATE","RENT_END_DATE","OWNER_NICKNAME","BORROWER_NICKNAME","STATUS","REVIEW_STATUS") VALUES(:1 ,:2 ,:3 ,:4 ,:5 ,:6 ,:7 ), Original SQL = INSERT INTO "RENTED_PRODUCT"("PRODUCT_NO","RENT_START_DATE","RENT_END_DATE","OWNER_NICKNAME","BORROWER_NICKNAME","STATUS","REVIEW_STATUS") VALUES(?,?,?,?,?,?,?), Error Message = ORA-01400: cannot insert NULL into ("TEST1"."RENTED_PRODUCT"."RENTED_PRODUCT_NO") at io.confluent.connect.jdbc.sink.JdbcSinkTask.getAllMessagesException(JdbcSinkTask.java:165) at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:111) ... 12 more [2023-12-07 10:20:17,979] ERROR [rentedproduct2-sink-connect|task-0] WorkerSinkTask{id=rentedproduct2-sink-connect-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:212) org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception. at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:618) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:336) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:237) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:206) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:181) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:834)
连接器配置
{ "name": "rentedproduct2-sink-connect", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "RENTED_PRODUCT", "connection.url": "jdbc:oracle:thin:@localhost:1521/xe", "connection.user": "test1", "connection.password": "test1", "insert.mode": "insert", "pk.mode": "none", "auto.create": "false", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable":"false", "value.converter.schemas.enable":"true" } }
Kafka消息格式
{ "schema": { "type": "struct", "fields": [ { "type": "int64", "optional": true, "field": "PRODUCT_NO" }, { "type": "string", "optional": true, "field": "RENT_START_DATE" }, { "type": "string", "optional": true, "field": "RENT_END_DATE" }, { "type": "string", "optional": true, "field": "OWNER_NICKNAME" }, { "type": "string", "optional": true, "field": "BORROWER_NICKNAME" }, { "type": "string", "optional": true, "field": "STATUS" }, { "type": "string", "optional": true, "field": "REVIEW_STATUS" } ], "optional": false, "name": "RENTED_PRODUCT" }, "payload": { "product_no": 1, "rent_start_date": "2023-12-12 11:00", "rent_end_date": "2023-12-15 12:00", "owner_nickname": "Owener", "borrower_nickname": "Borrower", "status": "Using", "review_status": "None" } }
Oracle建表语句(JPA生成)
create table rented_product (rented_product_no number(19,0) not null, borrower_nickname varchar2(255 char) not null, owner_nickname varchar2(255 char) not null, product_no number(19,0) not null, rent_end_date timestamp not null, rent_start_date timestamp not null, review_status varchar2(255 char), status varchar2(255 char), primary key (rented_product_no))
JPA实体部分代码
@Id @SequenceGenerator( name = "SEQ_GENERATOR", sequenceName = "RENTEDPRODUCT_SEQ", allocationSize = 1 ) @GeneratedValue(strategy = GenerationType.SEQUENCE, generator = "SEQ_GENERATOR") private Long rentedProductNo; @Column(nullable = false) private Long productNo; @Column(nullable = false) @JsonFormat(shape = JsonFormat.Shape.STRING, pattern = "yyyy-MM-dd HH:mm", timezone = "Asia/Seoul") private LocalDateTime rentStartDate; @Column(nullable = false) @JsonFormat(shape = JsonFormat.Shape.STRING, pattern = "yyyy-MM-dd HH:mm", timezone = "Asia/Seoul") private LocalDateTime rentEndDate; @Column(nullable = false) private String ownerNickname; @Column(nullable = false) private String borrowerNickname; @Enumerated(EnumType.STRING) private Status status; @Enumerated(EnumType.STRING) private ReviewStatus reviewStatus;
解决方案
原因分析
MySQL的自增主键会在INSERT语句未指定值时自动生成,但Oracle的序列不会自动作用于INSERT操作(除非字段设置默认值或通过触发器触发)。当前JPA的序列生成逻辑仅在应用层生效,而Kafka消息未携带主键字段,导致Sink Connector生成的INSERT语句不包含RENTED_PRODUCT_NO,Oracle因此抛出主键为空的错误。
解决方法
方法1:给Oracle主键字段设置默认值调用序列
修改表结构,让主键默认取序列的下一个值:
ALTER TABLE RENTED_PRODUCT MODIFY RENTED_PRODUCT_NO DEFAULT RENTEDPRODUCT_SEQ.NEXTVAL;
后续INSERT语句即使不指定该字段,Oracle会自动填充序列值。
方法2:自定义Connector的INSERT语句
在连接器配置中添加自定义插入语句,手动调用序列:
"custom.insert.statement": "INSERT INTO RENTED_PRODUCT(RENTED_PRODUCT_NO, PRODUCT_NO, RENT_START_DATE, RENT_END_DATE, OWNER_NICKNAME, BORROWER_NICKNAME, STATUS, REVIEW_STATUS) VALUES(RENTEDPRODUCT_SEQ.NEXTVAL, :PRODUCT_NO, :RENT_START_DATE, :RENT_END_DATE, :OWNER_NICKNAME, :BORROWER_NICKNAME, :STATUS, :REVIEW_STATUS)"
注意占位符需与Kafka消息中的字段名对应(Oracle不区分大小写,但建议保持统一)。
方法3:上游生产者携带主键值
修改上游应用,在发送Kafka消息前通过JPA或直接调用Oracle序列生成RENTED_PRODUCT_NO,并将该字段加入消息的schema和payload中,让Sink Connector直接写入值。
方法4:创建Oracle触发器自动生成主键
创建触发器,在INSERT操作时自动给主键赋值:
CREATE OR REPLACE TRIGGER TRG_RENTED_PRODUCT_NO BEFORE INSERT ON RENTED_PRODUCT FOR EACH ROW BEGIN SELECT RENTEDPRODUCT_SEQ.NEXTVAL INTO :NEW.RENTED_PRODUCT_NO FROM DUAL; END; /
此方法无需修改应用或Connector配置,触发器会自动处理主键生成。
内容的提问来源于stack exchange,提问作者khakyy

