This project sets up a Change Data Capture (CDC) pipeline using PostgreSQL, Kafka, Zookeeper, and Debezium. It demonstrates how to capture changes from a PostgreSQL database and stream them via Kafka topics.
- Docker & Docker Compose installed
- Node.js (for running the message listener script)
- (Optional) Postman or cURL for API calls
This will build and start the required containers:
- Kafka
- PostgreSQL
- Zookeeper
- Debezium connector
To build and start all containers, run:
docker-compose up --buildIf you want to use your local PostgreSQL database, run:
docker compose --profile local-db up -dAdd the connector configuration by sending the following POST request (using Postman or cURL):
URL: http://localhost:8083/connectors
Method: POST
Request Body (JSON):
{
"name": "postgres-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "postgres",
"database.password": "postgres",
"database.dbname": "testdb",
"database.server.name": "testdb-server",
"plugin.name": "pgoutput",
"slot.name": "debezium_slot_2",
"publication.name": "debezium_pub_3",
"table.include.list": "public.users",
"tombstones.on.delete": "false",
"key.converter.schemas.enable": "false",
"value.converter.schemas.enable": "false",
"topic.prefix": "testdb-server",
"snapshot.mode": "always"
}
}Access the Kafka UI at: http://localhost:8080/
A Node.js script (index.js) is included which listens to the Kafka topic and logs incoming messages.
Run it with:
node index.js- If messages are not appearing on the Kafka topic, try executing the SQL commands in
init.sqlmanually. - Delete and recreate the connector configuration if issues persist.
- Verify that Kafka and PostgreSQL containers are running and accessible.
- Check network ports (default Kafka: 9092, Debezium: 8083).
- PostgreSQL Logical Decoding Explanation
- Debezium PostgreSQL Connector Documentation
- Debezium/Postgres CDC Demo (YouTube)
- Ensure
wal_levelis set tologicalin PostgreSQL settings. - Create a user with
REPLICATIONprivilege:
SELECT rolname, rolsuper, rolreplication
FROM pg_roles
WHERE rolname = current_user;
-- If not assigned, run:
ALTER ROLE USERNAME WITH REPLICATION;