I am trying to figure out why I am receiving an "IDENTITY_INSERT" error when trying to use a JDBC sink connector to sink data into a SQL Server database off of a topic that is also written to by a JDBC source connector connected to the same SQL Server database.
The overall objective:
Currently there is a SQL Server database being used by the backend for storage in a traditional sense, and we are trying to transition to using Kafka for all of the same purposes, however the SQL Server database must remain for the time being as there are services that still rely upon it, and we have a requirement that all data that is on Kafka be mirrored in the SQL Server database.
What I am trying to accomplish:
I am trying to create a setup in which there are the following:
- One SQL Server database (all tables with identical primary key of "id" which auto-increments and is set by SQL Server)
- Kafka cluster, including Kafka connect with:
- JDBC Source connector to sync what is in the SQL Server table onto a kafka topic, lets call it AccountType for both the topic and the table
- JD Sink connector that subscribes to the same topic AccountType and sinks data into the same AccountType table in the SQL Server Database
The expected behavior is:
- if a legacy service writes/updates a record in SQL Server
- the source connector will pick up the change and write it to it's corresponding Kafka topic
- The sink connector will receive the message on that same topic, however, since the change originated in SQL Server and therefore has already been made from the perspective of the sink connector the sink connector will find a match on the primary key, see that there is no change to be made, and move along
- if a new service, designed to work with Kafka, updates a record and writes it to the correct topic:
- the JDBC sink connector will receive the message on the topic as an offset
- since the sink connector was configured with upsert mode, it will find a match on the primary key in the target database and update the corresponding record in the target database
- the source connector will then detect the change, triggering it to write the change to the corresponding topic
- My assumption at this point is that one of two things will happen, either:
- the source connector will not write to the topic since it would only be duplicating the last message or
- the source connector will write the duplicate message to the topic, however it will be ignored by the sink as there would be no resulting database record changes
This expected behavior aligns with everything I have found in the documentation, and is as best as I can tell in accordance with the JDBC sink deep dive guide found here: https://rmoff.net/2021/03/12/kafka-connect-jdbc-sink-deep-dive-working-with-primary-keys/[kafka-connect-jdbc-sink-deep-dive-working-with-primary-keys][1]
What is happening instead:
- with the Kafka cluster all started up, and the database empty, both connectors are successfully created
- A row is inserted into a table in the database using an external service
- the source connector successfully picks up the change and writes a record to the topic on Kafka (which has been split by the transforms such that the field representing the SQL Server table PK has been extracted and set as the message key, and removed from the value)
- (the problem) the sink connector then receives the message on that topic and...
...here's the issue, based on several videos and examples I could find, nothing should happen since that record is already up to date in the database, instead, however, it immediately attempts to write the entire message as-is, to the target table resulting in the following:
java.sql.BatchUpdateException: Cannot insert explicit value for identity column in table 'AccountType' when IDENTITY_INSERT is set to OFF.
Which makes sense, since the message from the topic has a primary key field in it, if it isn't enabled in the table then it shouldn't be allowed. Just for fun, I tried throwing in an additional transform to remove the id field before trying to write, and instead using another field in the table that has a "Unique" constraint in the configuration. When I repeated the steps this time it did not complain about writing the primary key, however it did still immediately try and insert the record which resulted in another error since it would violate the unique constraint, which again makes perfect sense.
Where I am stuck:
If all of the above makes sense can anyone tell me why it is automatically trying to insert despite being set to upsert?
Notes:
- all of this is being set up using docker containers provided by confluent for the confluent platform version 6.2.0
Source connector configuration:
{
"connection.url": "jdbc:sqlserver://mssql:1433;databaseName=REDACTED",
"connection.user":"REDACTED",
"connection.password":"REDACTED",
"connection.attempts": "3",
"connection.backoff.ms": "5000",
"table.whitelist": "AccountType",
"db.timezone": "UTC",
"name": "sql-server-source",
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"dialect.name": "SqlServerDatabaseDialect",
"config.action.reload": "restart",
"topic.creation.enable": "false",
"tasks.max": "1",
"mode": "timestamp+incrementing",
"incrementing.column.name": "id",
"timestamp.column.name": "created, updated",
"validate.non.null": true,
"key.converter": "org.apache.kafka.connect.converters.LongConverter",
"value.converter": "io.confluent.connect.json.JsonSchemaConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081",
"auto.register.schemas": "true",
"schema.registry.url": "http://schema-registry:8081",
"errors.log.include.messages": "true",
"transforms": "copyFieldToKey,extractKeyFromStruct,removeKeyFromValue",
"transforms.copyFieldToKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
"transforms.copyFieldToKey.fields": "id",
"transforms.extractKeyFromStruct.type":
"org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractKeyFromStruct.field": "id",
"transforms.removeKeyFromValue.type":
"org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.removeKeyFromValue.blacklist": "id",
"transforms": "copyFieldToKey,extractKeyFromStruct,removeKeyFromValue",
"transforms.copyFieldToKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
"transforms.copyFieldToKey.fields": "id",
"transforms.extractKeyFromStruct.type":
"org.apache.kafka.connect.transforms.ExtractField$Key",
"transforms.extractKeyFromStruct.field": "id",
"transforms.removeKeyFromValue.type":
"org.apache.kafka.connect.transforms.ReplaceField$Value",
"transforms.removeKeyFromValue.blacklist": "id",
}
Sink connector configuration:
{
"connection.url": "jdbc:sqlserver://mssql:1433;databaseName=REDACTED",
"connection.user":"REDACTED",
"connection.password":"REDACTED",
"connection.attempts": "3",
"connection.backoff.ms": "5000",
"table.name.format": "${topic}",
"db.timezone": "UTC",
"name": "sql-server-sink",
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"dialect.name": "SqlServerDatabaseDialect",
"auto.create": "false",
"auto.evolve": "false",
"tasks.max": "1",
"batch.size": "1000",
"topics": "AccountType",
"key.converter": "org.apache.kafka.connect.converters.LongConverter",
"value.converter": "io.confluent.connect.json.JsonSchemaConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081",
"insert.mode": "UPSERT",
"pk.mode": "record_key",
"pk.fields": "id",
}