Stream Data from PostgreSQL to Elasticsearch Without Writing Producer or Consumer Code.
In the previous articles in this Kafka series, I focused on Kafka’s delivery guarantees and what happens when messages move between producers, brokers, and consumers. In this article, I want to look at another important part of the Kafka ecosystem: Kafka Connect.
Suppose you already have data in PostgreSQL and want that data available in Kafka. From there, you also want the same data indexed in Elasticsearch.
One approach would be to build a Java application that reads PostgreSQL and produces records to Kafka. Then you could build another application that consumes those records and writes them to Elasticsearch. But for this kind of integration, we don’t necessarily need to write either application.
Kafka Connect is designed specifically for moving data between Kafka and external systems using reusable connectors. In the example from this article, we will build this pipeline: PostgreSQL → Kafka → Elasticsearch.
If you prefer to follow the practical demonstration as a video, you can watch it here.
What Is Kafka Connect?
Kafka Connect is part of the Apache Kafka ecosystem. Its purpose is to move data between Kafka and external systems reliably and at scale.

Instead of implementing a custom Kafka producer every time we need to bring external data into Kafka, or a custom consumer every time we want to send Kafka data somewhere else, Kafka Connect provides a framework where those integrations can be configured as connectors. There are two main types of connectors.
A source connector reads data from an external system and writes it into Kafka.
A sink connector does the opposite: it consumes data from Kafka and writes it to an external system.
For the example in this article:
- PostgreSQL is our source
- Kafka sits in the middle
- Elasticsearch is our destination

The Environment Used in the Demo
For this demonstration, everything runs locally with Docker Compose.
The Kafka environment is the same three-controller, three-broker KRaft setup I used in the previous videos, with a small change to the volumes. The Kafka brokers use three partitions by default, replication factor 3, and min.insync.replicas=2.
In addition to Kafka, I run PostgreSQL and Elasticsearch. Redis and Adminer are also present in the environment, although they are not used for this particular example. The final component is our Kafka Connect cluster.
For this experiment I use two Kafka Connect workers:
connect-1:
image: confluentinc/cp-kafka-connect:7.9.2
ports:
- 8081:8083
connect-2:
image: confluentinc/cp-kafka-connect:7.9.2
ports:
- 8082:8083
Both workers belong to the same Connect cluster:
CONNECT_GROUP_ID: payment-connect-cluster
and both communicate with all three Kafka brokers:
CONNECT_BOOTSTRAP_SERVERS: broker-1:9092,broker-2:9092,broker-3:9092
They also share the same internal topics for configuration, offsets, and status:
CONNECT_CONFIG_STORAGE_TOPIC: connect-configs
CONNECT_OFFSET_STORAGE_TOPIC: connect-offsets
CONNECT_STATUS_STORAGE_TOPIC: connect-status
The worker configuration also defines the plugin directory:
CONNECT_PLUGIN_PATH: '/usr/share/java,/plugins'
CONNECT_PLUGIN_DISCOVERY: 'only_scan'
and maps /plugins to a directory on the host containing the connector plugins. The second worker uses the same Connect group, Kafka brokers, internal topics, converters, and plugin location. This gives us a two-worker Kafka Connect cluster rather than a single standalone Connect process.
Installing the Connector Plugins
Kafka Connect provides the runtime, but it still needs connector plugins that know how to communicate with the systems we want to integrate.
For this demonstration, we need two:
JDBC Source Connector for PostgreSQL and Elasticsearch Sink Connector for Elasticsearch.
As I explain in the video, I install the plugins using the Confluent Hub client because the installation also handles the connector dependencies.
The downloaded plugins are placed in the external plugin directory mounted into the Connect containers.
When the workers start, they scan that directory through:
CONNECT_PLUGIN_PATH: '/usr/share/java,/plugins'
This makes the installed connector implementations available to Kafka Connect.
Preparing PostgreSQL, Kafka, and Elasticsearch
Before creating any connectors, it is useful to look at the starting state.
The PostgreSQL database contains two simple tables:
payment
cart
At the beginning of the demonstration, payment contains one record and cart contains two records.
Kafka already has corresponding topics:
paymentcart
but both topics are initially empty. Elasticsearch also has corresponding indices: payment and cart and these are empty as well.
This gives us a very clear starting point: PostgreSQL contains the data. Kafka does not. Elasticsearch does not. Now we can watch the data move through each stage.
Creating the JDBC Source Connector
We begin with the source side. The Kafka Connect JDBC Source Connector reads relational database tables and produces their rows into Kafka.
An important point from the demonstration is that the JDBC source connector is table-oriented rather than relationship-aware. It does not automatically navigate relationships or perform joins between tables.
With the table-based configuration used here, the mapping is straightforward:
PostgreSQL payment table → Kafka payment topic and PostgreSQL cart table → Kafka cart topic.
The source connector configuration used in the project is:
{
"name": "payment-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"connection.url": "jdbc:postgresql://postgres:5432/payment",
"connection.user": "payment",
"connection.password": "payment",
"table.whitelist": "payment,cart",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "",
"poll.interval.ms": "5000",
"tasks.max": "1"
}
}
The most important part for understanding how this particular connector detects new data is:
"mode": "incrementing",
"incrementing.column.name": "id"
The connector uses the id column to determine which rows are new.
Because:
"topic.prefix": ""
does not add anything before the table name, the resulting Kafka topics retain the table names.
Registering a Connector Through the REST API
Defining the JSON configuration is not enough by itself. Kafka Connect exposes a REST API that we use to create and manage connectors. The video calls this out early: to make Kafka Connect perform the integration, we register the connector through its REST API.
After registering payment-source, Kafka Connect creates the connector with one task. Almost immediately, we can see the source worker reading the existing rows from PostgreSQL.
Remember the state before creating the connector:
PostgreSQL: contains records
Kafka: empty
Elasticsearch: empty
Now we consume the cart topic from Kafka. The PostgreSQL cart records are there. Then we inspect payment. The payment record is there as well.
At this point our pipeline has reached only this far: PostgreSQL → JDBC Source Connector → Kafka We have not created the sink connector yet.
If we query Elasticsearch, both destination indices are still empty. That separation is useful because it demonstrates exactly what each connector is responsible for.
What About Updates and More Complex Queries?
The example uses:
"mode": "incrementing"
which works well for detecting newly inserted rows based on an increasing ID. But the video also discusses several other possibilities. If you need to detect updates, the JDBC connector can use a timestamp column:
mode=timestamp
timestamp.column.name=updated_at
Another possibility is combining the two:
mode=timestamp+incrementing
In that case, the connector uses both an incrementing column such as id and a timestamp column to identify changes. The demonstration also mentions using a custom SQL query rather than simply providing a table whitelist.
That becomes useful when the data you want to produce into Kafka does not correspond directly to one database table for example, when you want a query involving multiple tables to produce records into a single topic.
Creating the Elasticsearch Sink Connector
We now have the PostgreSQL data in Kafka. The next step is to move it from Kafka to Elasticsearch. This is the responsibility of our sink connector. The Elasticsearch sink connector configuration subscribes to both topics:
"topics": "payment,cart"
and connects to:
"connection.url": "http://elasticsearch:9200"
The project configuration is:
{
"name": "payment-sink",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"topics": "payment,cart",
"connection.url": "http://elasticsearch:9200",
"key.ignore": "true",
"schema.ignore": "true",
"behavior.on.malformed.documents": "warn",
"write.method": "UPSERT",
"tasks.max": "1"
}
}
Unlike the JDBC source connector, which is reading PostgreSQL, this connector is consuming Kafka records. In this example the mapping is again straightforward:
Kafka payment topic → Elasticsearch payment index
Kafka cart topic → Elasticsearch cart index
The Complete Kafka Connect Pipeline

Verifying the Data in Elasticsearch
Before we created the sink connector, the Elasticsearch indices were empty even though Kafka already contained the PostgreSQL records.
After registering the sink connector, we query Elasticsearch again. This time, the payment index contains the payment data. The cart index contains the cart records as well.
The full pipeline is now working:
PostgreSQL → Kafka → Elasticsearch
Kafka Connect continuously handles the integration. The source connector polls PostgreSQL for new records and writes them into Kafka, while the sink connector consumes the Kafka topics and writes the records into Elasticsearch. And we achieved this without implementing a custom Java producer or consumer. That is the real value Kafka Connect provides for this kind of integration.
Managing Kafka Connect Through Its REST API
Creating connectors is only one part of the REST API.
For example, we can list the currently registered connectors:
curl "$CONNECT_URL/connectors"
We can check the status of a connector:
curl "$CONNECT_URL/connectors/payment-source/status"
We can pause it:
curl -X PUT "$CONNECT_URL/connectors/payment-source/pause"
resume it:
curl -X PUT "$CONNECT_URL/connectors/payment-source/resume"
or restart it:
curl -X POST "$CONNECT_URL/connectors/payment-source/restart"
And when a connector is no longer needed, it can be deleted through the REST API as well.
These endpoints become particularly useful when operating Connect beyond a simple local demonstration because connectors and their tasks have their own lifecycle that needs to be observed and managed.
Source Connectors vs. Sink Connectors
A source connector brings external data into Kafka.
In our example: PostgreSQL → Kafka
A sink connector takes Kafka data somewhere else.
In our example: Kafka → Elasticsearch
Kafka sits between them, which means the source and destination integrations remain separated.
The JDBC source connector does not need to know that Elasticsearch exists.
The Elasticsearch sink connector does not need to know that the records originally came from PostgreSQL. Both sides integrate through Kafka. That separation is one of the useful properties of this architecture.
No Custom Producer or Consumer Code
The title of the video says “No Coding!”, and this experiment demonstrates what that means in practice.
We did not implement something like:
KafkaProducer<String, Payment> producer;
to extract the database rows.
Nor did we implement:
KafkaConsumer<String, Payment> consumer;
to read Kafka and index the records in Elasticsearch. Instead, we described the integration through configuration.
The JDBC Source Connector knows how to communicate with PostgreSQL. Kafka Connect knows how to execute and manage the connector.
The Elasticsearch Sink Connector knows how to consume Kafka records and communicate with Elasticsearch.
Our responsibility becomes configuring and operating the pipeline rather than writing the plumbing ourselves.
When Kafka Connect Fits
This example is intentionally simple, but it demonstrates the problem Kafka Connect is designed to solve.
When data already exists in systems such as databases and needs to flow into Kafka—or Kafka records need to be delivered to another external system—a supported connector can remove a significant amount of integration code.
The same basic architecture is not limited to PostgreSQL and Elasticsearch. Kafka Connect is intended to bridge Kafka with external systems through reusable connectors.
But it is also important to choose the connector based on the actual problem. The JDBC source demonstrated here polls database tables. When the requirement becomes full change-data capture—including more sophisticated handling of inserts, updates and deletes.
Final Thoughts
Kafka Connect gives us a standardized way to move data between Kafka and external systems without building a custom producer or consumer for every integration.
In this hands-on example, we started with records stored in PostgreSQL while both Kafka and Elasticsearch were empty.
We configured a JDBC source connector, and the PostgreSQL records appeared in the payment and cart Kafka topics.
At that stage Elasticsearch was still empty.
We then created the Elasticsearch sink connector. It consumed those Kafka topics and wrote the records into the corresponding Elasticsearch indices.
The result was a complete streaming integration:
PostgreSQL → JDBC Source Connector → Apache Kafka → Elasticsearch Sink Connector → Elasticsearch
All driven by Kafka Connect configuration and its REST API rather than application code.
If you want to see the entire setup—including the Docker environment, two-worker Kafka Connect cluster, plugin installation, PostgreSQL tables, connector creation, Kafka topics, Elasticsearch indices, and REST API operations—you can follow the complete hands-on demonstration here.