JDBC Sink Without a Schema
First written December 2022, last updated September 2026.
Not all Kafka integration tools are the same. Some integration systems only produce JSON data without a schema. The JDBC Sink Connector requires a schema. Here are steps showcasing a low-code option to push events into a relational database when the source data is schema-less JSON.
Real-Time Toolkit
This demo is available in the postgres-sink folder of the dev-local-demos project. It uses applications available through containers in dev-local. All you need is a Kafka Cluster with the Confluent Schema Registry, two ksqlDB queries per topic, and a JDBC Sink Connector running on a Connect cluster.
Challenges
Relational databases require a schema. The JDBC Sink connector relies on the JDBC API and database-specific drivers to write data from a Kafka topic into a table on the database. These drivers need to know the data types of the fields being written.
Steps
Two ksqlDB queries are enough to give schema-less JSON a schema the JDBC Sink can use.
Work with streams, not tables
If the goal is to stream the events into the relational database table rather than persist the data within ksqlDB for stream processing, use streams (not tables) within ksqlDB. This allows Kafka topic retention to be only as long as the applications need to process the events.
Read Data into a Stream
The goal of the first statement is to capture the data in ksqlDB with no transformation. This query needs to align with how the data is coming from the source system. This example shows a client that publishes a complete message as JSON with the order_id also the key of the message.
create or replace stream ORDERS ( "order_id" varchar key, "user_id" varchar, "store_id" varchar, "quantity" bigint, "amt" decimal(4,2), "ts" string) with (kafka_topic='ORDERS', value_format='json', key_format='kafka');Transform and publish with schema
Specify a key and value format that has a schema; for demonstration, Avro is used. Also, renaming fields to align with the database table’s schema will minimize the number of single message transformations in the connect configuration.
create stream ORDERS_WITH_SCHEMA with(KEY_FORMAT='avro', VALUE_FORMAT='avro')asselect "order_id" as ORDER_ID, "user_id" as USER_ID, "store_id" as STORE_ID, "quantity" as QUANTITY, "amt" as AMT, PARSE_TIMESTAMP("ts", 'yyyy-MM-dd HH:mm:ss', 'UTC') as TSfrom ORDERS;While this query is good, it can be better. By ensuring the key is itself a record (struct), it will simplify the JDBC Sink Connector configuration and allow for a single instance to handle multiple topics.
create or replace stream ORDERS_WITH_SCHEMA with(KEY_FORMAT='avro', VALUE_FORMAT='avro')asSELECT STRUCT(ORDER_ID:=`order_id`) as PK, "user_id" as USER_ID, "store_id" as STORE_ID, "quantity" as QUANTITY, "amt" as AMT, PARSE_TIMESTAMP("ts", 'yyyy-MM-dd HH:mm:ss', 'UTC') as TSFROM ORDERSPARTITION BY STRUCT(ORDER_ID:=`order_id`);When the key is a primitive type, the schema specification captures only the data type and drops the field name. By making it a record, a schema is associated with the key, which is then accessible by the JDBC Sink Connector through the Schema Registry. This eliminates the need to set the pk.fields in the connector’s configuration.
Foreign Keys
With a topic-per-table setup, disabling foreign-key constraints in the destination database is required, since order guarantees aren't achievable between tables. This is a big architecture discussion to settle before designing a migration of data to a relational database.
If this isn't possible, a single topic for parent/child events is an option, where the message key is the parent’s primary key and the JDBC sink connector is configured with the primary key pulled from the value. This is a fair amount more development and configuration. If the source topic is an aggregate with parents and their children together, then a custom consumer could be the right solution. These are discussions outside this article and tutorial.
Connector Configuration
When it comes to a specific Apache Kafka connector, please read the configuration documentation closely, as there are numerous differences between connector configurations. With the JDBC Sink connector, a few specific ones should be called out.
Dialect
The dialect tells the connector the type of database, which determines how the SQL is constructed.
When it comes to Postgres, pay close attention to the name. PostgresDatabaseDialect and PostgresSqlDatabaseDialect are incorrect.
"dialect.name": "PostgreSqlDatabaseDialect"Connection secrets
For connection information, put these in a secret and reference them with the provider.
The file provider, FileConfigProvider, ships with the Apache Kafka distribution and is easy to enable.
"connection.url": "${file:/etc/kafka-connect/secrets/postgres.properties:CONNECTION_URL}","connection.user" : "${file:/etc/kafka-connect/secrets/postgres.properties:CONNECTION_USER}","connection.password" : "${file:/etc/kafka-connect/secrets/postgres.properties:CONNECTION_PASSWORD}"Putting the connection.url value in a secret makes it easier to reuse configuration across environments.
Quoting identifiers
Do not change the default of quote.sql.identifiers from always to never, as invalid JDBC statements may result.
"quote.sql.identifiers": "always"In the referenced demonstration, changing this to never results in invalid syntax because of the AMT field.
Insert mode
Understand the behavior of the insert mode. If you have an idempotent system and the database isn’t generating sequence numbers as part of the insert, upsert will be the typical setting.
"insert.mode": "upsert"This setting is the easiest to use and understand in a streaming platform. The syntax of the merge (upsert) SQL statements can cause unexpected database performance issues. Enabling trace logging will show the prepared statements in the logs, and doing an explain-plan analysis on these queries can help address performance concerns.
pk.mode and pk.fields
Setting pk.mode to record.key is the easiest way to deploy a single connector for multiple tables with different primary key field names. If the key schema is a primitive, pk.fields is required. However, if the key is a structure, then this can (and should) be omitted, since all the fields in the key’s structure are used as the primary key.
partition.assignment.strategy
If you are configuring multiple topics for one connector, you need to understand partition assignment.
The default strategy is RangeAssignor. This means that the same partitions for all topics are shared within the same consumer instance. Not only can you only have as many tasks as the partition count of the topic with the most partitions, but you will also end up with uneven distributions if the partition counts aren't the same.
In the provided demonstration, ORDERS_WITH_SCHEMA has 4 partitions and USERS_WITH_SCHEMA has 2 partitions, which yields uneven workers and limits the number of tasks to 4, even if tasks.max=6.
Range Assignor
kafka-consumer-groups --bootstrap-server <broker>:9092 --describe --group connect-postgres
| GROUP | TOPIC | PARTITION | CLIENT-ID |
|---|---|---|---|
| connect-postgres | ORDERS_WITH_SCHEMA | 0 | connector-consumer-postgres-0 |
| connect-postgres | USERS_WITH_SCHEMA | 0 | connector-consumer-postgres-0 |
| connect-postgres | ORDERS_WITH_SCHEMA | 1 | connector-consumer-postgres-1 |
| connect-postgres | USERS_WITH_SCHEMA | 1 | connector-consumer-postgres-1 |
| connect-postgres | ORDERS_WITH_SCHEMA | 2 | connector-consumer-postgres-2 |
| connect-postgres | ORDERS_WITH_SCHEMA | 3 | connector-consumer-postgres-3 |
Only 4 working tasks are created (unique client-ids), since the maximum partition count of the topics is 4. Because USERS has two partitions, two workers get two partitions each and two workers get one.
Round Robin Assignor
"consumer.override.partition.assignment.strategy": "org.apache.kafka.clients.consumer.RoundRobinAssignor",kafka-consumer-groups --bootstrap-server <broker>:9092 --describe --group connect-postgres
tasks.max = 6
| GROUP | TOPIC | PARTITION | CLIENT-ID |
|---|---|---|---|
| connect-postgres | ORDERS_WITH_SCHEMA | 0 | connector-consumer-postgres-0 |
| connect-postgres | ORDERS_WITH_SCHEMA | 1 | connector-consumer-postgres-1 |
| connect-postgres | ORDERS_WITH_SCHEMA | 2 | connector-consumer-postgres-2 |
| connect-postgres | ORDERS_WITH_SCHEMA | 3 | connector-consumer-postgres-3 |
| connect-postgres | USERS_WITH_SCHEMA | 0 | connector-consumer-postgres-4 |
| connect-postgres | USERS_WITH_SCHEMA | 1 | connector-consumer-postgres-5 |
The maximum number of possible tasks is the total number of partitions across all topics. If the setting is lower, such as tasks.max=5, partitions are assigned in a round-robin fashion (hence the name).
tasks.max = 5
| GROUP | TOPIC | PARTITION | CLIENT-ID |
|---|---|---|---|
| connect-postgres | ORDERS_WITH_SCHEMA | 0 | connector-consumer-postgres-0 |
| connect-postgres | USERS_WITH_SCHEMA | 1 | connector-consumer-postgres-0 |
| connect-postgres | ORDERS_WITH_SCHEMA | 1 | connector-consumer-postgres-1 |
| connect-postgres | ORDERS_WITH_SCHEMA | 2 | connector-consumer-postgres-2 |
| connect-postgres | ORDERS_WITH_SCHEMA | 3 | connector-consumer-postgres-3 |
| connect-postgres | USERS_WITH_SCHEMA | 0 | connector-consumer-postgres-4 |
Takeaways
- You can't always change how a source system publishes. A stream processor in the middle lets you fix it without asking.
- Existing Single Message Transforms (SMTs) don't bridge the schema-less to schema transformation. Writing a custom SMT is an option, but wouldn't be available to use in a Confluent Cloud managed Connector.
- A simple Kafka Streams application would achieve the same result. While a Kafka Streams solution gives more flexibility, it increases development and operational effort.
- While the goal of streaming systems isn't a pass-through from a source system into RDBMS tables, many enterprise systems need some level of this functionality.
- Foreign-key constraints lead to challenges; disabling them for an “eventually consistent” database may not be an option. A more advanced solution is needed if they exist in the source tables.
- Check out the dev-local project and the demo in
postgres-sinkin dev-local-demos, and try it out yourself.
Working on something like this?
Start a Conversation