Skip to content

Latest commit

 

History

History
479 lines (351 loc) · 17.2 KB

File metadata and controls

479 lines (351 loc) · 17.2 KB

Part 3: Data Engineering: From PostgreSQL to ClickHouse via Debezium CDC

Overview

In Part 2 of the lab, we built a containerized microservices architecture where a producer published orders to Kafka and two consumers processed them — one sending notifications, and the other persisting to PostgreSQL.

In Part 3 of the lab, we will build a data pipeline on top of that foundation. Using Debezium Change Data Capture (CDC), you will stream every change made to the PostgreSQL orders table into Kafka, transform the data, and load it into a ClickHouse data warehouse in real time.

The architecture is shown below.


System Architecture

System Architecture


New Concepts in This Lab

Change Data Capture (CDC)

In Part 2, the Inventory Service wrote orders to PostgreSQL directly. But what if you have an existing application that writes to a database and you cannot modify its code? For example, a business that outsourced the development of one of its business applications and the developers did not make the code open-source. However, the business has full access to the database used by the business application. CDC solves such cases.

Instead of changing the application, CDC taps into the database's internal change log — the Write-Ahead Log (WAL) — and streams every INSERT, UPDATE, and DELETE as an event to Kafka. The application never knows CDC is running.

This is how large-scale data pipelines are built in practice: you capture changes at the database level, not at the application level.

Write Ahead Log (WAL)

The Write-Ahead Log is PostgreSQL's durability mechanism. Every change to data — every INSERT, UPDATE, DELETE — is written to the WAL before it is written to the actual data files.

This serves two purposes:

  1. Crash recovery — if the server dies mid-write, the WAL is replayed on restart
  2. Replication — other systems can consume the WAL to replicate or, in this case, stream changes.

The wal_level setting controls how much information is written into each WAL record.

Level What It Records Use Case
minimal Bare minimum for crash recovery Standalone, no replication
replica Enough for physical/streaming replication Read replicas
logical Full row-level change detail, including old values Change Data Capture, logical replication

Enabling wal_level=logical increases the WAL size — PostgreSQL ends up writing more data per transaction. This has a significant storage and I/O cost in a production environment.

The traditional safeguard is to set max_slot_wal_keep_size in PostgreSQL >=13. This caps how much WAL is retained per slot and drops the slot rather than filling the disk. If Debezium goes offline and PostgreSQL starts accumulating unread WAL logs beyond 6GB, PostgreSQL will protect itself by invalidating the slot rather than filling your disk and crashing.

Debezium

Debezium is a CDC platform that runs as a Kafka Connect plugin. It connects to PostgreSQL as a replication client, reads the WAL, and publishes every change as a structured JSON event to a Kafka topic.

Each Debezium event is set up to have the following structure:

{
  "payload": {
    "before": null,
    "after": {
      "order_id": "abc-123",
      "client_fname": "Jeff",
      "item": "Managu",
      "order_quantity": 3,
      "received_at": 1714000000000000
    },
    "op": "c"
  }
}

The op field tells you what kind of change occurred:

op Meaning
c Create — a new row was INSERTED
r Snapshot read — existing row read when the connector was started
u Update — an existing row was UPDATED
d Delete — a row was DELETED

ClickHouse

ClickHouse is a columnar database built for analytical queries.

Feature PostgreSQL (row store) ClickHouse (column store)
Storage layout One row stored together One column stored together
Best for INSERT / SELECT single rows Aggregations over millions of rows
Use case Transactional systems (OLTP) Analytics and reporting (OLAP)
Example query "Fetch order #1234" "Total revenue by item this month"

In a production data platform, PostgreSQL holds the live operational data and ClickHouse holds the historical analytical data for dashboards and reports.

It is desirable to have the data in the data warehouse or data lake or data lakehouse be as close to real-time as possible. Analogy: When crossing a road, you want to look both ways and have the most up-to-date information about oncoming traffic. If what you see is 1 hour old, then it is like looking both ways but only seeing where the cars were yesterday — and that is not very helpful for making decisions in the present.

The Transformation Step

Raw Debezium events reflect the source database schema as it is. A transformation step exists in the pipeline to:

  • Rename fields to match the warehouse naming convention
  • Compute derived fields that would be expensive to recalculate at query time
  • Convert data types (e.g., microsecond timestamps to datetime objects)
  • Filter events that should not enter the warehouse, e.g., removing Personally Identifiable Information (PII) to comply with privacy regulations in the Kenya Data Protection Act.
  • Enrich records with metadata that the source system does not store

etc.

The following transformations are applied in the transformer.py script in this lab:

Transformation Source value Warehouse value
Field rename client_fname customer_name
Computed field order_quantity is_bulk_order (if order_quantity > 5 then true (1) else false (0))
Timestamp Microseconds since epoch (integer) Python datetime object with a timezone (Nairobi time in this case; UTC+3)
Added audit field — processed_at (pipeline time)
Added audit field op code operation (readable label)
Filter op = "d" (DELETE) Skipped entirely

Step-by-Step Instructions

Prerequisites

Ensure you have completed Part 2 and understand the producer, consumer, and broker concepts first because Part 3 of the lab builds directly on top of the concepts in Part 2.


Navigate into the Part 3 directory (3_data_engineering) first. All the commands below assume that you are inside the 3_data_engineering directory.

cd 3_data_engineering

Step 1: Run Unit Tests

Run the unit tests:

cd producer/
pytest -v -s test_producer_order.py
cd ../consumer-notification/
pytest -v -s test_consumer_order_notification.py
cd ../consumer-inventory/
pytest -v -s test_consumer_order_inventory.py
cd ..

Step 2: Create the .env File

Create a file called .env inside the 3_data_engineering directory and paste the contents of the .env.example file inside it.

Set the values of the environment variables as discussed in class.

Step 3: Set Up Directories and Start the Stack

# Create the required volume directories
chmod u+x project_setup.sh
sed -i.bak 's/\r$//' project_setup.sh
./project_setup.sh
# Build all images and start all services
docker compose -f docker-compose.yaml up --build \
  --scale producer=1 \
  --scale consumer-notification=1 \
  --scale consumer-inventory=1

This will start:

  • 3 Kafka brokers (kafka1, kafka2, kafka3)
  • 1 PostgreSQL container
  • 1 Kafka Connect container with Debezium
  • 1 ClickHouse container
  • 1 producer service
  • 1 notification consumer service
  • 1 inventory consumer service
  • 1 transformer service

Wait until all containers are running and healthy before proceeding. You can check their status with:

docker-compose ps

All the containers should be in a healthy status after a few minutes. If any container is not healthy, check its logs to debug the issue.


Step 4: Verify the Stack is Healthy

Check that Kafka brokers are up:

docker exec kafka1 kafka-topics \
  --bootstrap-server localhost:9092 \
  --list

Check that PostgreSQL is receiving orders:

docker exec -it postgres psql -U lab_user -d lab_db \
  -c "SELECT * FROM orders ORDER BY received_at DESC LIMIT 5;"

Check that Kafka Connect is ready:

curl http://localhost:8083/

You should see a JSON response showing the Kafka Connect version. If you receive a connection error, wait 30 more seconds and try again.


Step 5: Register the Debezium Connector

This is the step that activates CDC. You are basically telling Debezium which database and table to monitor.

Important: An explanation of cURL is available here and the documentation of the connector configuration being sent as JSON data via HTTP POST using cURL is available here.

NOTE: This should be executed from inside the 3_data_engineering/ directory:

chmod u+x kafka-connect/register-connector.sh
sed -i 's/\r$//' kafka-connect/register-connector.sh
./kafka-connect/register-connector.sh

Expected output:

Script directory = ./kafka-connect
Waiting for Kafka Connect to be ready...
Kafka Connect is ready.

Registering Debezium PostgreSQL connector...
✅ Connector registered successfully (HTTP 201).

You can verify the connector status with:
  curl http://localhost:8083/connectors/orders-postgres-connector/status

What the connector configuration does (connector-config.json):

Field Value Purpose
connector.class PostgresConnector Uses the Debezium PostgreSQL connector
database.hostname postgres Docker service name of the PostgreSQL container
topic.prefix dbserver1 Prefix for all Kafka topics this connector creates
table.include.list public.orders Only monitor this table (schema.table format)
plugin.name pgoutput Uses PostgreSQL's built-in logical decoding plugin (no extra installation needed)
slot.name debezium_slot The replication slot Debezium creates in PostgreSQL to track its position in the WAL
publication.autocleanup.on.connector.stop true Cleans up the PostgreSQL replication slot when the connector is stopped
snapshot.mode initial Read all existing rows first, then switch to live CDC

Verify the connector is running:

curl http://localhost:8083/connectors/orders-postgres-connector/status

The state field should show RUNNING.


Step 6: Verify the Debezium Topic Exists

Once the connector is registered, Debezium creates a new Kafka topic named dbserver1.public.orders and immediately begins publishing snapshot events for all existing rows in the PostgreSQL orders table.

# List all topics — you should now see dbserver1.public.orders
docker exec kafka1 kafka-topics \
  --bootstrap-server localhost:9092 \
  --list

# Inspect the CDC topic — notice that it has 3 partitions and 3 replicas
docker exec kafka1 kafka-topics \
  --bootstrap-server localhost:9092 \
  --describe --topic dbserver1.public.orders

# Consume raw events from the CDC topic to see the Debezium format
docker exec kafka1 kafka-console-consumer \
  --bootstrap-server localhost:9092 \
  --topic dbserver1.public.orders \
  --from-beginning \
  --max-messages 3

Scheme through the raw JSON output. Identify the payload.op, payload.before, and payload.after fields. This is the exact input that transformer.py receives and processes.


Step 7: Verify Data is Arriving in ClickHouse

The transformer service is already running as a Docker container. Check its logs to see it consuming CDC events and writing to ClickHouse:

docker-compose logs transformer

Then query ClickHouse directly to see the transformed data:

# Connect to ClickHouse using the CLI
docker exec -it clickhouse clickhouse-client

# Once inside the CLI, run these queries:

-- Count total orders in the warehouse
SELECT count() FROM orders;

-- View all orders
SELECT * FROM orders ORDER BY processed_at DESC LIMIT 10;

-- See the effect of the is_bulk_order field
SELECT
    is_bulk_order,
    count()        AS order_count,
    sum(order_quantity) AS total_units
FROM orders
GROUP BY is_bulk_order;

-- Compare received_at and processed_at to measure pipeline latency
SELECT
    order_id,
    received_at,
    processed_at,
    dateDiff('second', received_at, processed_at) AS latency_seconds
FROM orders
ORDER BY processed_at DESC
LIMIT 10;

-- Exit the CLI
exit;

Step 8: Observe the Live Pipeline

At this point the full pipeline is running continuously. Open connections to PostgreSQL and ClickHouse using DataGrip or any SQL client of your choice to observe the data in both databases in real time.

DataGrip Output

You can also use the terminal to watch the data flow if you are using Linux or MacOS.

If you are using Linux:

sudo apt install procps

If you are using MacOS:

brew install watch

Then open three terminal windows and run the following simultaneously to observe the end-to-end flow in real time.

Terminal 1 — Watch new orders arrive in PostgreSQL:

watch -n 2 "docker exec postgres psql -U lab_user -d lab_db \
  -c 'SELECT order_id, item, order_quantity, received_at FROM orders \
  ORDER BY received_at DESC LIMIT 5;'"

Terminal 2 — Watch the transformer processing CDC events:

docker-compose logs -f transformer

Terminal 3 — Watch ClickHouse receiving transformed records:

watch -n 2 "docker exec clickhouse clickhouse-client \
  --query 'SELECT order_id, customer_name, item, is_bulk_order, operation \
  FROM orders ORDER BY processed_at DESC LIMIT 5'"

You should see rows appearing in PostgreSQL and then — within seconds — appearing in ClickHouse with the applied transformations. The customer_name column will contain the renamed value and is_bulk_order will be set automatically based on the quantity.


Step 9: Observe CDC Operations (INSERT, UPDATE, DELETE)

The producer only inserts new orders. To observe UPDATE and DELETE events flowing through the pipeline, run these commands manually.

Trigger an UPDATE:

docker exec -it postgres psql -U lab_user -d lab_db -c "
UPDATE orders
SET order_quantity = 99
WHERE order_id = (SELECT order_id FROM orders LIMIT 1);
"

Watch the transformer logs — you will see an event with operation: UPDATE. Then query ClickHouse to see the updated record.

Trigger a DELETE:

docker exec -it postgres psql -U lab_user -d lab_db -c "
DELETE FROM orders
WHERE order_id = (SELECT order_id FROM orders LIMIT 1);
"

Watch the transformer logs — you will see Skipping DELETE event. The record will remain in ClickHouse because the transformer filters out DELETE operations by design. This demonstrates an intentional architectural decision: the warehouse retains historical data even when the source database removes it.


Step 10: Tear Down (Project Cleanup)

# Stop all services AND delete all stored data
docker-compose down -v
chmod u+x project_cleanup.sh
sed -i.bak 's/\r$//' project_cleanup.sh
./project_cleanup.sh