Kafka-connect-jdbc: kafka jdbc sink connectors to mysql incorrect detect database as present

Created on 27 Jun 2019  路  1Comment  路  Source: confluentinc/kafka-connect-jdbc

my kafka bootstrap.servers=10.0.10.7:9092,10.0.200.15:9092,10.0.200.11:9092
i start connect at machine 10.0.10.7 with ./bin/connect-distributed -daemon ./etc/kafka/connect-distributed.properties
this is my connect-distributed.properties:

bootstrap.servers=10.0.10.7:9092,10.0.200.15:9092,10.0.200.11:9092
group.id=connect-cluster
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true
offset.storage.topic=connect-offsets
offset.storage.replication.factor=1
config.storage.topic=connect-configs
config.storage.replication.factor=1
status.storage.topic=connect-status
status.storage.replication.factor=1
offset.flush.interval.ms=10000
plugin.path=/opt/third/confluent-5.2.1/share/java

jdbc driver path is /opt/third/confluent-5.2.1/share/java/kafka-connect-jdbc/mysql-connector-java-8.0.16.jar
add connector with POST http://10.0.10.7:8083/connectors
this is my request body:

{
    "name":"test_mysql_sink",
    "config":{
        "connector.class":"io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max":1,
        "topics":"wltest_mysql_sink",
        "connection.url":"jdbc:mysql://127.0.0.1:3306/testmysink",
        "connection.user":"root",
        "connection.password":"123456",
        "auto.create":true
    }
}

i produce data at machine 10.0.10.7 with: ./bin/kafka-console-producer --broker-list 10.0.10.7:9092,10.0.200.15:9092,10.0.200.11:9092 --topic wltest_mysql_sink
my json data is
{"schema":{"type":"struct","fields":[{"type":"int32","optional":true,"field":"c1"},{"type":"string","optional":true,"field":"c2"},{"type":"int64","optional":false,"name":"org.apache.kafka.connect.data.Timestamp","version":1,"field":"create_ts"},{"type":"int64","optional":false,"name":"org.apache.kafka.connect.data.Timestamp","version":1,"field":"update_ts"}],"optional":false,"name":"wltest_mysql_sink"},"payload":{"c1":10000,"c2":"bar","create_ts":1501834166000,"update_ts":1501834166000}}

my error is:

[2019-06-27 17:16:54,883] INFO Attempting to open connection #1 to MySql (io.confluent.connect.jdbc.util.CachedConnectionProvider:87)
[2019-06-27 17:16:55,131] INFO JdbcDbWriter Connected (io.confluent.connect.jdbc.sink.JdbcDbWriter:49)
[2019-06-27 17:16:55,144] INFO Checking MySql dialect for existence of table "wltest_mysql_sink" (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:511)
[2019-06-27 17:16:55,174] INFO Using MySql dialect table "wltest_mysql_sink" absent (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:519)
[2019-06-27 17:16:55,177] INFO Creating table with sql: CREATE TABLE `wltest_mysql_sink` (
`update_ts` DATETIME(3) NOT NULL,
`create_ts` DATETIME(3) NOT NULL,
`c1` INT NULL,
`c2` VARCHAR(256) NULL) (io.confluent.connect.jdbc.sink.DbStructure:92)
[2019-06-27 17:16:56,008] INFO Checking MySql dialect for existence of table "wltest_mysql_sink" (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:511)
[2019-06-27 17:16:56,016] INFO Using MySql dialect table "wltest_mysql_sink" present (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:519)
[2019-06-27 17:16:56,027] WARN Write of 1 records failed, remainingRetries=1 (io.confluent.connect.jdbc.sink.JdbcSinkTask:76)
java.sql.SQLSyntaxErrorException: Table 'IP.wltest_mysql_sink' doesn't exist
    at com.mysql.cj.jdbc.exceptions.SQLError.createSQLException(SQLError.java:120)
    at com.mysql.cj.jdbc.exceptions.SQLError.createSQLException(SQLError.java:97)
    at com.mysql.cj.jdbc.exceptions.SQLExceptionsMapping.translateException(SQLExceptionsMapping.java:122)
    at com.mysql.cj.jdbc.StatementImpl.executeQuery(StatementImpl.java:1218)
    at com.mysql.cj.jdbc.DatabaseMetaData$7.forEach(DatabaseMetaData.java:2980)
    at com.mysql.cj.jdbc.DatabaseMetaData$7.forEach(DatabaseMetaData.java:2968)
    at com.mysql.cj.jdbc.IterateBlock.doForAll(IterateBlock.java:56)
    at com.mysql.cj.jdbc.DatabaseMetaData.getPrimaryKeys(DatabaseMetaData.java:3021)
    at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.primaryKeyColumns(GenericDatabaseDialect.java:717)
    at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.describeColumns(GenericDatabaseDialect.java:554)
    at io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.describeTable(GenericDatabaseDialect.java:753)
    at io.confluent.connect.jdbc.util.TableDefinitions.get(TableDefinitions.java:62)
    at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:112)
    at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
    at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:85)
    at io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:66)
    at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:74)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:538)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:321)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    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)
[2019-06-27 17:16:56,029] INFO Closing connection #1 to MySql (io.confluent.connect.jdbc.util.CachedConnectionProvider:113)
[2019-06-27 17:16:56,029] INFO Initializing writer using SQL dialect: MySqlDatabaseDialect (io.confluent.connect.jdbc.sink.JdbcSinkTask:57)
[2019-06-27 17:16:56,030] ERROR WorkerSinkTask{id=test_mysql_sink-0} RetriableException from SinkTask: (org.apache.kafka.connect.runtime.WorkerSinkTask:551)
org.apache.kafka.connect.errors.RetriableException: java.sql.SQLException: java.sql.SQLSyntaxErrorException: Table 'IP.wltest_mysql_sink' doesn't exist

    at io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:93)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:538)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:321)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:224)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:192)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:175)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:219)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    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: java.sql.SQLException: java.sql.SQLSyntaxErrorException: Table 'IP.wltest_mysql_sink' doesn't exist

    ... 12 more
[2019-06-27 17:16:57,186] INFO Attempting to open connection #1 to MySql (io.confluent.connect.jdbc.util.CachedConnectionProvider:87)
[2019-06-27 17:16:57,193] INFO JdbcDbWriter Connected (io.confluent.connect.jdbc.sink.JdbcDbWriter:49)
[2019-06-27 17:16:57,196] INFO Checking MySql dialect for existence of table "wltest_mysql_sink" (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:511)
[2019-06-27 17:16:57,203] INFO Using MySql dialect table "wltest_mysql_sink" present (io.confluent.connect.jdbc.dialect.MySqlDatabaseDialect:519)
[2019-06-27 17:16:57,204] WARN Write of 1 records failed, remainingRetries=0 (io.confluent.connect.jdbc.sink.JdbcSinkTask:76)
java.sql.SQLSyntaxErrorException: Table 'IP.wltest_mysql_sink' doesn't exist

It just create table wltest_mysql_sink in db testmysink success.But insert data err, IP database is another in mysql nothing to do with this, why it scan other databases , rather than to use db(testmysink) in connection url.

Most helpful comment

in connection.url should add nullCatalogMeansCurrent=true :
"connection.url":"jdbc:mysql://127.0.0.1:3306/testmysink?nullCatalogMeansCurrent=true",
because MySQL Connector/J 8.0 changes nullCatalogMeansCurrent to false which cause DatabaseMetaData.getTables will return tables in all db

>All comments

in connection.url should add nullCatalogMeansCurrent=true :
"connection.url":"jdbc:mysql://127.0.0.1:3306/testmysink?nullCatalogMeansCurrent=true",
because MySQL Connector/J 8.0 changes nullCatalogMeansCurrent to false which cause DatabaseMetaData.getTables will return tables in all db

Was this page helpful?
0 / 5 - 0 ratings