Kafka-connect-jdbc: is it possible to avoid quoting names of tables and colums is postgres

Created on 26 Jul 2018  路  5Comments  路  Source: confluentinc/kafka-connect-jdbc

I have an Oracle db as source and PosgreSQL as sink.
Table and column names in oracle are in upper case, to preserve this kafka-connect-jdbc uses quotes in ddl.
Is it possible to change this behavior and have a non quoted columns and table names in PosgreSQL?

Most helpful comment

This is still not working for me using v5.3.1. When not using the quoting option

curl -s -X POST -H  "Content-Type:application/json" \
    http://localhost:8083/connectors/ \
    -d '
{
    "name": "sysda_avro_postgres",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "1",
        "topics": "SYSDA_AVRO",
        "connection.url": "jdbc:postgresql://localhost:5432/kafkasink",
        "connection.user": "postgres",
        "connection.password": "somepw",
        "auto.create": "true",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        "key.converter.schemas.enable": "false",
        "value.converter.schemas.enable": "true",
        "value.converter.schema.registry.url": "http://localhost:8081"
    }
}
'

tables are created and populated with data correctly, but table and columns names are quoted (which is usually not desired).

When adding the option `"quote.sql.identifiers": "never"`` to prevent quoting of identifiers, it fails with

[2019-10-07 19:30:34,028] INFO Attempting to open connection #1 to PostgreSql (io.confluent.connect.jdbc.util.CachedConnectionProvider:87)
[2019-10-07 19:30:34,041] INFO JdbcDbWriter Connected (io.confluent.connect.jdbc.sink.JdbcDbWriter:49)
[2019-10-07 19:30:34,052] INFO Checking PostgreSql dialect for existence of table "SYSDA_AVRO" (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:511)
[2019-10-07 19:30:34,059] INFO Using PostgreSql dialect table "SYSDA_AVRO" absent (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:519)
[2019-10-07 19:30:34,060] INFO Creating table with sql: CREATE TABLE SYSDA_AVRO (
ORRS BIGINT NULL,
PR BIGINT NULL,
VER TEXT NULL,
TIME TEXT NULL,
IPPORT TEXT NULL,
PM BIGINT NULL) (io.confluent.connect.jdbc.sink.DbStructure:92)
[2019-10-07 19:30:34,066] INFO Checking PostgreSql dialect for existence of table "SYSDA_AVRO" (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:511)
[2019-10-07 19:30:34,068] INFO Using PostgreSql dialect table "SYSDA_AVRO" absent (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:519)
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. (org.apache.kafka.connect.runtime.WorkerSinkTask:558)
java.lang.NullPointerException
        at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:124)
        at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
        at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:121)
        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:177)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)
        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)
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:179)
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:560)
        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:177)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)
        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)
Caused by: java.lang.NullPointerException
        at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:124)
        at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
        at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:121)
        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)
        ... 10 more
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:180)
[2019-10-07 19:30:34,069] INFO Stopping task (io.confluent.connect.jdbc.sink.JdbcSinkTask:105)
[2019-10-07 19:30:34,069] INFO Closing connection #1 to PostgreSql (io.confluent.connect.jdbc.util.CachedConnectionProvider:113)
[2019-10-07 19:30:34,070] INFO [Consumer clientId=connector-consumer-sysda_avro_postgres-0, groupId=connect-sysda_avro_postgres] Member connector-consumer-sysda_avro_postgres-0-cb9c42d5-ab81-48ce-a672-aee484f6b2f9 sending LeaveGroup request to coordinator dddocker02:9092 (id: 2147483647 rack: null) (org.apache.kafka.clients.consumer.internals.AbstractCoordinator:879)
[2019-10-07 19:30:34,072] INFO Publish thread interrupted for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:285)
[2019-10-07 19:30:34,076] INFO [Producer clientId=confluent.monitoring.interceptor.connector-consumer-sysda_avro_postgres-0] Cluster ID: S-S5top9RGG4uxhaPWXtdw (org.apache.kafka.clients.Metadata:261)
[2019-10-07 19:30:34,077] INFO Publishing Monitoring Metrics stopped for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:297)
[2019-10-07 19:30:34,077] INFO [Producer clientId=confluent.monitoring.interceptor.connector-consumer-sysda_avro_postgres-0] Closing the Kafka producer with timeoutMillis = 9223372036854775807 ms. (org.apache.kafka.clients.producer.KafkaProducer:1153)
[2019-10-07 19:30:34,083] INFO Closed monitoring interceptor for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:320)

Initially, it seems to have picked up the quoting option, because in the logged schema there is not quoting (which is the case when using "always"). However, it then fails with some error. Feel welcome if you need further information, or in case you know any workaround.

All 5 comments

It is possible to bypass these quotes with #572, which was recently made and has not yet been released, though it will be in the 5.0.2, 5.1.1, and 5.2.0 releases whenever they occur. In the meantime, you can build the connector locally and try the new feature to disable quoting.

I'll mark this as closed, but feel free to reopen if you try the feature and it does not work for you.

This is still not working for me using v5.3.1. When not using the quoting option

curl -s -X POST -H  "Content-Type:application/json" \
    http://localhost:8083/connectors/ \
    -d '
{
    "name": "sysda_avro_postgres",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "tasks.max": "1",
        "topics": "SYSDA_AVRO",
        "connection.url": "jdbc:postgresql://localhost:5432/kafkasink",
        "connection.user": "postgres",
        "connection.password": "somepw",
        "auto.create": "true",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        "key.converter.schemas.enable": "false",
        "value.converter.schemas.enable": "true",
        "value.converter.schema.registry.url": "http://localhost:8081"
    }
}
'

tables are created and populated with data correctly, but table and columns names are quoted (which is usually not desired).

When adding the option `"quote.sql.identifiers": "never"`` to prevent quoting of identifiers, it fails with

[2019-10-07 19:30:34,028] INFO Attempting to open connection #1 to PostgreSql (io.confluent.connect.jdbc.util.CachedConnectionProvider:87)
[2019-10-07 19:30:34,041] INFO JdbcDbWriter Connected (io.confluent.connect.jdbc.sink.JdbcDbWriter:49)
[2019-10-07 19:30:34,052] INFO Checking PostgreSql dialect for existence of table "SYSDA_AVRO" (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:511)
[2019-10-07 19:30:34,059] INFO Using PostgreSql dialect table "SYSDA_AVRO" absent (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:519)
[2019-10-07 19:30:34,060] INFO Creating table with sql: CREATE TABLE SYSDA_AVRO (
ORRS BIGINT NULL,
PR BIGINT NULL,
VER TEXT NULL,
TIME TEXT NULL,
IPPORT TEXT NULL,
PM BIGINT NULL) (io.confluent.connect.jdbc.sink.DbStructure:92)
[2019-10-07 19:30:34,066] INFO Checking PostgreSql dialect for existence of table "SYSDA_AVRO" (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:511)
[2019-10-07 19:30:34,068] INFO Using PostgreSql dialect table "SYSDA_AVRO" absent (io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect:519)
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. (org.apache.kafka.connect.runtime.WorkerSinkTask:558)
java.lang.NullPointerException
        at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:124)
        at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
        at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:121)
        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:177)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)
        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)
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task threw an uncaught and unrecoverable exception (org.apache.kafka.connect.runtime.WorkerTask:179)
org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
        at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:560)
        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:177)
        at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)
        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)
Caused by: java.lang.NullPointerException
        at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:124)
        at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
        at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:121)
        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)
        ... 10 more
[2019-10-07 19:30:34,069] ERROR WorkerSinkTask{id=sysda_avro_postgres-0} Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask:180)
[2019-10-07 19:30:34,069] INFO Stopping task (io.confluent.connect.jdbc.sink.JdbcSinkTask:105)
[2019-10-07 19:30:34,069] INFO Closing connection #1 to PostgreSql (io.confluent.connect.jdbc.util.CachedConnectionProvider:113)
[2019-10-07 19:30:34,070] INFO [Consumer clientId=connector-consumer-sysda_avro_postgres-0, groupId=connect-sysda_avro_postgres] Member connector-consumer-sysda_avro_postgres-0-cb9c42d5-ab81-48ce-a672-aee484f6b2f9 sending LeaveGroup request to coordinator dddocker02:9092 (id: 2147483647 rack: null) (org.apache.kafka.clients.consumer.internals.AbstractCoordinator:879)
[2019-10-07 19:30:34,072] INFO Publish thread interrupted for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:285)
[2019-10-07 19:30:34,076] INFO [Producer clientId=confluent.monitoring.interceptor.connector-consumer-sysda_avro_postgres-0] Cluster ID: S-S5top9RGG4uxhaPWXtdw (org.apache.kafka.clients.Metadata:261)
[2019-10-07 19:30:34,077] INFO Publishing Monitoring Metrics stopped for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:297)
[2019-10-07 19:30:34,077] INFO [Producer clientId=confluent.monitoring.interceptor.connector-consumer-sysda_avro_postgres-0] Closing the Kafka producer with timeoutMillis = 9223372036854775807 ms. (org.apache.kafka.clients.producer.KafkaProducer:1153)
[2019-10-07 19:30:34,083] INFO Closed monitoring interceptor for client_id=connector-consumer-sysda_avro_postgres-0 client_type=CONSUMER session= cluster=S-S5top9RGG4uxhaPWXtdw group=connect-sysda_avro_postgres (io.confluent.monitoring.clients.interceptor.MonitoringInterceptor:320)

Initially, it seems to have picked up the quoting option, because in the logged schema there is not quoting (which is the case when using "always"). However, it then fails with some error. Feel welcome if you need further information, or in case you know any workaround.

I'm using Confluent Platform version 5.4.0 and i'm facing the same problem.
With the option "quote.sql.identifiers":"always", it works but i need quote to request my table =>
select "myField" from "myTable";
If i change this option to "quote.sql.identifiers":"never", the connector raises this exception:
_org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:561)
at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:322)
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:177)
at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)
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.lang.NullPointerException
at io.confluent.connect.jdbc.sink.DbStructure.amendIfNecessary(DbStructure.java:124)
at io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:74)
at io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:121)
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:539)
... 10 more_

I found a workaround. If the destination table already exist and was correctly created (without quotes), then the connector works correctly only if using the option "quote.sql.identifiers":"never".

Regards

BUT it doesn't work if the destination table doesn't exist and if using the option "auto.create": true to create the table. With these condition, the table is created using quotes.

Was this page helpful?
0 / 5 - 0 ratings