DTStack / DTStack/chunjun

flinkx 执行批量同步mysql-mysql, update模式:执行一直都是报错

Open
#1,017 4 comments 0 reactions 0 assignees View on GitHub
bug
Dominant language
Java
Stars
4.1k
Forks
1.7k
PR merge metrics
No merged PRs in 30d

Description

BUG信息:
JdbcOutputFormat [Flink_Job] writeRecord error: when converting field[0] in Row(+I(18,xulei,22,wuhan))
at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357)
at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1915)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.quiesceTimeServiceAndCloseOperator(StreamOperatorWrapper.java:168)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:131)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.close(StreamOperatorWrapper.java:135)
at org.apache.flink.streaming.runtime.tasks.OperatorChain.closeOperators(OperatorChain.java:439)
at org.apache.flink.streaming.runtime.tasks.StreamTask.afterInvoke(StreamTask.java:627)
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:589)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:755)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:570)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.RuntimeException: java.lang.IllegalArgumentException: WritingRecordError: error writing record [2] exceed limit [0]
+I(18,xulei,22,wuhan)
com.dtstack.flinkx.throwable.WriteRecordException:
JdbcOutputFormat [Flink_Job] writeRecord error: when converting field[0] in Row(+I(18,xulei,22,wuhan))
com.mysql.jdbc.exceptions.jdbc4.MySQLIntegrityConstraintViolationException: Duplicate entry '18' for key 'PRIMARY'
at com.dtstack.flinkx.connector.jdbc.sink.JdbcOutputFormat.processWriteException(JdbcOutputFormat.java:342)
at com.dtstack.flinkx.connector.jdbc.sink.JdbcOutputFormat.writeSingleRecordInternal(JdbcOutputFormat.java:181)
at com.dtstack.flinkx.sink.format.BaseRichOutputFormat.writeSingleRecord(BaseRichOutputFormat.java:465)
at java.util.ArrayList.forEach(ArrayList.java:1249)
at com.dtstack.flinkx.sink.format.BaseRichOutputFormat.writeRecordInternal(BaseRichOutputFormat.java:485)
at com.dtstack.flinkx.sink.format.BaseRichOutputFormat.lambda$initTimingSubmitTask$0(BaseRichOutputFormat.java:438)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180)
at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Caused by: com.mysql.jdbc.exceptions.jdbc4.MySQLIntegrityConstraintViolationException: Duplicate entry '18' for key 'PRIMARY'
at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
at sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
at sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
at com.mysql.jdbc.Util.handleNewInstance(Util.java:425)
at com.mysql.jdbc.Util.getInstance(Util.java:408)
at com.mysql.jdbc.SQLError.createSQLException(SQLError.java:936)
at com.mysql.jdbc.MysqlIO.checkErrorPacket(MysqlIO.java:3976)
at com.mysql.jdbc.MysqlIO.checkErrorPacket(MysqlIO.java:3912)
at com.mysql.jdbc.MysqlIO.sendCommand(MysqlIO.java:2530)
at com.mysql.jdbc.MysqlIO.sqlQueryDirect(MysqlIO.java:2683)
at com.mysql.jdbc.ConnectionImpl.execSQL(ConnectionImpl.java:2486)
at com.mysql.jdbc.PreparedStatement.executeInternal(PreparedStatement.java:1858)
at com.mysql.jdbc.PreparedStatement.execute(PreparedStatement.java:1197)
at com.dtstack.flinkx.connector.jdbc.statement.FieldNamedPreparedStatementImpl.execute(FieldNamedPreparedStatementImpl.java:76)
at com.dtstack.flinkx.connector.jdbc.sink.JdbcOutputFormat.writeSingleRecordInternal(JdbcOutputFormat.java:175)
... 11 more

JdbcOutputFormat [Flink_Job] writeRecord error: when converting field[0] in Row(+I(18,xulei,22,wuhan))
at com.dtstack.flinkx.sink.format.BaseRichOutputFormat.close(BaseRichOutputFormat.java:332)
at com.dtstack.flinkx.sink.DtOutputFormatSinkFunction.close(DtOutputFormatSinkFunction.java:127)
at org.apache.flink.api.common.functions.util.FunctionUtils.closeFunction(FunctionUtils.java:41)
at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.close(AbstractUdfStreamOperator.java:109)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$closeOperator$5(StreamOperatorWrapper.java:213)
at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.runThrowing(StreamTaskActionExecutor.java:93)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.closeOperator(StreamOperatorWrapper.java:210)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.lambda$deferCloseOperatorToMailbox$3(StreamOperatorWrapper.java:185)
at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.runThrowing(StreamTaskActionExecutor.java:93)
at org.apache.flink.streaming.runtime.tasks.mailbox.Mail.run(Mail.java:90)
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxExecutorImpl.tryYield(MailboxExecutorImpl.java:97)
at org.apache.flink.streaming.runtime.tasks.StreamOperatorWrapper.quiesceTimeServiceAndCloseOperator(StreamOperatorWrapper.java:162)
... 8 more
Caused by: java.lang.IllegalArgumentException: WritingRecordError: error writing record [2] exceed limit [0]
+I(18,xulei,22,wuhan)
com.dtstack.flinkx.throwable.WriteRecordException:
JdbcOutputFormat [Flink_Job] writeRecord error: when converting field[0] in Row(+I(18,xulei,22,wuhan))
com.mysql.jdbc.exceptions.jdbc4.MySQLIntegrityConstraintViolationException: Duplicate entry '18' for key 'PRIMARY'

执行JSON:
{
"job": {
"content": [
{
"reader": {
"parameter": {
"password": "123456",
"dataSourceId": 38,
"column": [
{
"precision": 10,
"name": "id",
"columnDisplaySize": 10,
"type": "INT"
},
{
"precision": 20,
"name": "name",
"columnDisplaySize": 20,
"type": "VARCHAR"
},
{
"precision": 10,
"name": "age",
"columnDisplaySize": 10,
"type": "INT"
},
{
"precision": 20,
"name": "address",
"columnDisplaySize": 20,
"type": "VARCHAR"
}
],
"connection": [
{
"jdbcUrl": [
"jdbc:mysql://172.18.8.113:3306/test_fjf"
],
"table": [
"mysqlreader"
]
}
],
"splitPk": "id",
"username": "root"
},
"name": "mysqlreader"
},
"writer": {
"parameter": {
"password": "123456",
"dataSourceId": 38,
"updateKey": [
"id"
],
"column": [
{
"precision": 10,
"name": "id",
"columnDisplaySize": 10,
"type": "INT"
},
{
"precision": 20,
"name": "name",
"columnDisplaySize": 20,
"type": "VARCHAR"
},
{
"precision": 10,
"name": "age",
"columnDisplaySize": 10,
"type": "INT"
},
{
"precision": 20,
"name": "address",
"columnDisplaySize": 20,
"type": "VARCHAR"
}
],
"connection": [
{
"jdbcUrl": "jdbc:mysql://172.18.8.113:3306/test_fjf",
"table": [
"mysqlwriter"
]
}
],
"writeMode": "update",
"username": "root"
},
"name": "mysqlwriter"
}
}
],
"setting": {

"speed": {
"bytes": 0,
"channel": 1
}
}
}
}

执行表:
![1656660167(1)](https://user-images.githubusercontent.com/34857750/176845161-ebd469a3-2a44-4231-a851-22749a09d658.png)
![1656660194(1)](https://user-images.githubusercontent.com/34857750/176845238-8497b4e2-cfab-498e-bf7e-1d0a771a05d5.png)

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.