Most databases change all day: users sign up, orders are placed, and each order moves from one status to the next. Change data capture (CDC) turns those changes into a stream of events that other systems can read as they happen. In this post, a simulated online shop writes to PostgreSQL in real time, Debezium streams every change to Kafka, and a sink connector saves the changes as files in object storage. Everything runs on your own machine.

The shop’s data model follows the theLook eCommerce dataset. The source code is in benchtop/ecommerce-cdc.

What You Will Build

  1. Run the simulation, which fills the tables and keeps changing them.
  2. Deploy the connectors: Debezium streams every insert and update to Kafka, and the S3 sink saves them as files.
  3. Look at the changes in Kafka and in SeaweedFS.
  4. Change the simulation while it runs, and see the change in the stream.

To start again at any point, python -m ecommerce.stores.cleanup removes everything this project has created and keeps the services running. Clean Up describes what it removes.

Architecture

The data moves in one direction: a simulation writes to PostgreSQL, Debezium reads the changes into Kafka topics, and the S3 sink writes the topics to files in SeaweedFS.

  • Change data capture means reading a database’s changes as a stream of events, instead of querying its tables again and again. PostgreSQL records every change in its write-ahead log (WAL) before it applies it, and CDC reads that log.
  • Kafka stores streams of events in topics. A topic keeps its events in order, and any number of readers can read it.
  • Kafka Connect runs connectors. A source connector copies data into Kafka, and a sink connector copies data out of it. You deploy a connector by sending its settings, as JSON, to Connect’s REST API.
  • Debezium is a source connector for CDC. It first copies every existing row, which is called a snapshot. It then streams each new change from the WAL, and writes one topic per table.
  • A replication slot is PostgreSQL’s record of how far a reader has got in the WAL. PostgreSQL keeps the part of the WAL the slot has not read yet, so Debezium can stop and carry on without losing a change.
  • A publication names the tables whose changes PostgreSQL sends. odctl’s PostgreSQL has a publication, cdc_pub, that covers every table in the cdc schema.
  • Aiven’s S3 sink is a sink connector. It reads the topics and writes their events to files in SeaweedFS, an object store with the same API as Amazon S3.
  • dynamic-des runs the simulation, a discrete-event model of the shop. dynamic-des is a Python library for simulations that stream their output.

The simulation writes six tables in PostgreSQL’s cdc schema:

TableRowsChanges
dist_centers10 distribution centresinserted once
products260 productsinserted once
usersregistered usersinserted, and updated when a user moves address
ordersordersinserted as Processing, then updated to Shipped, Delivered, Cancelled or Returned
order_itemsthe products in each orderinserted, and updated with their order’s status
eventspage views in web sessions, including anonymous visitorsinserted only

Setup

You need Docker (Docker Desktop, OrbStack or Docker Engine), uv and Python 3.13. Clone the repository and run every command from the project folder:

1git clone https://github.com/jaehyeon-kim/benchtop.git
2cd benchtop/ecommerce-cdc
3
4uv venv                             # create .venv
5source .venv/bin/activate           # activate it, in each new shell
6uv pip install -r requirements.txt
7
8odctl up postgres kafka-lite storage

odctl is a command line tool that starts a local data stack with Docker Compose. Its profiles start these services:

  • postgres: PostgreSQL, already set up for CDC. It runs with wal_level=logical, which makes the WAL hold enough detail for CDC. It also has the schema cdc with the publication cdc_pub.
  • kafka-lite: one Kafka broker, Kafka Connect with the Debezium and Aiven S3 connectors, Karapace and Kafka UI. Karapace is a schema registry: it stores the schema of each topic’s messages, so a message carries only a short schema id.
  • storage: SeaweedFS.

The web UIs:

  • Kafka UI: http://127.0.0.1:8086, for the topics, their messages and the connectors
  • SeaweedFS file browser: http://127.0.0.1:8889

Simulation Model

A discrete-event simulation (DES) moves a clock from one event to the next, such as a visitor arriving or an order being packed. Nothing happens between events, so the model only has to say what each event does and how long it is until the next one. Here the clock runs at the speed of real time.

The shop has three kinds of process. A process is a function that runs through simulated time and waits between its steps:

  • Visitor: visitors arrive at random, about one a second. Each views the home page and one to four products, and spends a few seconds on each page. Six in ten are registered users. Three in ten buy: they view the cart, sign up first if they are new, and place an order of one to four products.
  • Order: each order waits for one of the warehouse’s pickers. Three pickers pack orders, one at a time each. When orders arrive faster than they are packed, they queue. An order that waits longer than its customer’s patience, about two minutes, is cancelled. A packed order is shipped, delivered about a minute later, and one in ten is returned.
  • Move: now and then a registered user moves to a new address.

dynamic-des provides each part of this model:

Partdynamic-des featureHere
random arrivalsadd_arrival, arrival_loopvisitors, and address moves
time spent on a stepadd_servicea page view, packing, transit, patience and the time until a return
limited capacityadd_resourcethe pickers, which make orders queue
chancesadd_variablea visitor is registered, a visitor buys, an order is returned
a process of its ownspawneach visit and each order
live changesthe registry and KafkaIngressany of the above, while it runs
writesPostgresEgressone per table, upserting on id

Every parameter is in ecommerce/core/config.py. Each one is added to dynamic-des’s registry, a store of named parameters that the model reads each time it draws a value:

 1SERVICES = {  # mean and standard deviation, in seconds, lognormal
 2    "page_view": (4.0, 2.0),  # time on a page before the next one
 3    # how long an order waits for a picker before it is cancelled
 4    "patience": (120.0, 60.0),
 5    "pick": (5.0, 2.0),  # a picker packs an order
 6    "transit": (60.0, 20.0),  # a shipped order reaches the customer
 7    "return_after": (60.0, 30.0),  # a returned order comes back after delivery
 8}
 9CHANCES = {  # shares of visitors or orders, from 0 to 1
10    "returning": 0.6,  # a visitor is a registered user
11    "buy": 0.3,  # a visitor buys at the end of the visit
12    "return": 0.1,  # a delivered order is returned
13}
14PICKERS, MAX_PICKERS = 3, 10  # warehouse pickers packing orders at once

The order process is in ecommerce/simulation/run.py. It asks for a picker and a patience timeout at once, and picker | _wait(ctx, "patience") resumes at whichever comes first. If the timeout comes first, the order is cancelled. Each status change is published with the same id as before:

 1def fulfil(
 2    ctx: SimulationContext, state: Shop, order: Order, items: list[OrderItem]
 3) -> Process:
 4    """One order: packed by a picker or cancelled, then shipped, delivered and maybe returned."""
 5    with ctx.get_resource("pickers").request() as picker:
 6        got = yield picker | _wait(ctx, "patience")
 7        if picker not in got:
 8            _publish(ctx, *shop.advance(order, items, "Cancelled", _now(ctx, state)))
 9            return
10        yield _wait(ctx, "pick")
11    order, items = shop.advance(order, items, "Shipped", _now(ctx, state))
12    _publish(ctx, order, items)
13    yield _wait(ctx, "transit")
14    order, items = shop.advance(order, items, "Delivered", _now(ctx, state))
15    _publish(ctx, order, items)
16    if _chance(ctx, "return"):
17        yield _wait(ctx, "return_after")
18        _publish(ctx, *shop.advance(order, items, "Returned", _now(ctx, state)))

The run function adds one PostgreSQL writer per table, each upserting on id. An upsert inserts a row, or updates it when a row with the same key exists. So a changed order reaches PostgreSQL as an update, and Debezium sees it as one. The same function adds the Kafka reader for live changes:

 1def run(minutes: float | None = None, seed: int | None = None) -> None:
 2    """
 3    Runs the shop in real time, writing every row to its table and reading live changes.
 4
 5    Args:
 6        minutes (float, optional): How long to run. None runs until interrupted.
 7        seed (int, optional): The seed for every random choice.
 8    """
 9    postgres.create_tables()
10    KafkaAdminConnector(KAFKA).create_topics([{"name": CONTROL_TOPIC}])
11    app = build(datetime.now(UTC), seed)
12    app.add_ingress(KafkaIngress(topic=CONTROL_TOPIC, bootstrap_servers=KAFKA))
13    app.with_batching(batch_size=4, flush_interval=1.0)  # a few rows a second per table
14    for table in TABLES.values():
15        app.add_egress(
16            PostgresEgress(DSN, table_name=table, upsert_keys=KEY),
17            when=lambda r, t=table: r["stream_type"] == "event" and r["key"] == t,
18        )
19    app.run(until=minutes * 60 if minutes else None)

Times are stored as ISO 8601 text in UTC. dynamic-des publishes each row as JSON-friendly values, so a time reaches PostgreSQL as a string, and a timestamp column does not accept a string.

Step 1: Run the Simulation

1python -m ecommerce.simulation.run --minutes 16     # or leave out --minutes, and stop with Ctrl + C

It creates the six tables if they are missing, and writes the centres, the products and 100 users. It then runs in real time. With the default parameters it places about 20 orders a minute, and each order changes three or four times over the next few minutes. Its log shows the Kafka reader and one PostgreSQL writer per table:

2026-09-30 22:39:25,564 dynamic_des.core.context Building SimulationContext for 'ecommerce'...
2026-09-30 22:39:25,565 dynamic_des.core.context Simulation engine started.
2026-09-30 22:39:25,599 dynamic_des.connectors.ingress.kafka Connected to Kafka Ingress topic: ecommerce-control
2026-09-30 22:39:25,966 dynamic_des.connectors.egress.postgres PostgresEgress connected to dist_centers
2026-09-30 22:39:25,966 dynamic_des.connectors.egress.postgres PostgresEgress connected to products
2026-09-30 22:39:25,966 dynamic_des.connectors.egress.postgres PostgresEgress connected to events
2026-09-30 22:39:25,966 dynamic_des.connectors.egress.postgres PostgresEgress connected to order_items
2026-09-30 22:39:25,966 dynamic_des.connectors.egress.postgres PostgresEgress connected to orders
2026-09-30 22:39:25,967 dynamic_des.connectors.egress.postgres PostgresEgress connected to users

Leave it running, and use a second terminal for the next steps.

Step 2: Deploy the Connectors

1python -m ecommerce.cdc.connectors

It sends two JSON files to Kafka Connect. Connect creates each connector, or updates it if it exists. The source connector’s settings are in ecommerce/cdc/source.json:

 1{
 2  "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
 3  "tasks.max": "1",
 4  "database.hostname": "postgres",
 5  "database.port": "5432",
 6  "database.user": "user",
 7  "database.password": "password",
 8  "database.dbname": "odctl",
 9  "plugin.name": "pgoutput",
10  "publication.name": "cdc_pub",
11  "publication.autocreate.mode": "disabled",
12  "slot.name": "ecommerce_cdc",
13  "table.include.list": "cdc.users,cdc.products,cdc.dist_centers,cdc.orders,cdc.order_items,cdc.events",
14  "topic.prefix": "ecommerce",
15  "snapshot.mode": "initial",
16  "key.converter": "io.confluent.connect.avro.AvroConverter",
17  "key.converter.schema.registry.url": "http://karapace:8081",
18  "value.converter": "io.confluent.connect.avro.AvroConverter",
19  "value.converter.schema.registry.url": "http://karapace:8081"
20}

These settings tell Debezium to:

  • read the database odctl with PostgreSQL’s built-in pgoutput decoder;
  • use the existing publication cdc_pub, and never create one (publication.autocreate.mode is disabled);
  • keep its place in the WAL in its own replication slot, ecommerce_cdc;
  • capture only this project’s six tables, and name the topics ecommerce.cdc.<table>;
  • take a snapshot of the existing rows first (snapshot.mode is initial);
  • write Avro, a compact binary format, with each table’s key and value schemas registered in Karapace as the subjects ecommerce.cdc.<table>-key and ecommerce.cdc.<table>-value.

The sink connector’s settings are in ecommerce/cdc/s3-sink.json:

 1{
 2  "connector.class": "io.aiven.kafka.connect.s3.AivenKafkaConnectS3SinkConnector",
 3  "tasks.max": "1",
 4  "topics.regex": "ecommerce\\.cdc\\..*",
 5  "consumer.override.metadata.max.age.ms": "30000",
 6  "aws.access.key.id": "user",
 7  "aws.secret.access.key": "password",
 8  "aws.s3.endpoint": "http://seaweed:8333",
 9  "aws.s3.region": "us-east-1",
10  "aws.s3.bucket.name": "odctl-dev",
11  "file.name.template": "ecommerce-cdc/{{topic}}/{{partition}}-{{start_offset}}.jsonl",
12  "file.compression.type": "none",
13  "format.output.type": "jsonl",
14  "format.output.fields": "key,value,offset,timestamp",
15  "format.output.fields.value.encoding": "none",
16  "key.converter": "io.confluent.connect.avro.AvroConverter",
17  "key.converter.schema.registry.url": "http://karapace:8081",
18  "value.converter": "io.confluent.connect.avro.AvroConverter",
19  "value.converter.schema.registry.url": "http://karapace:8081"
20}

It reads every topic whose name matches ecommerce.cdc.*, and decodes each message with its schema from Karapace. It writes JSON lines files to the bucket odctl-dev, under ecommerce-cdc/<topic>/, with one change event on each line. The sink writes a file each time Connect commits its offsets, which record how far the sink has read. That happens about once a minute. The orders topic is created only when the first order is placed, after the sink has started. The sink subscribes by a topic pattern, and a Kafka consumer looks for new topics only when it refreshes its list of topics, every five minutes by default. consumer.override.metadata.max.age.ms makes the sink’s consumer refresh it every 30 seconds instead. In a test run, the first orders file appeared after about a minute.

Each connector’s status should show RUNNING for the connector and its task:

1curl -s http://127.0.0.1:8083/connectors/ecommerce-cdc-source/status
2curl -s http://127.0.0.1:8083/connectors/ecommerce-cdc-s3/status

Kafka UI lists both connectors under Kafka Connect.

Kafka UI listing the Debezium source and the S3 sink, both running
Both connectors in Kafka UI

Step 3: Look at the Changes

Debezium creates one topic per table. The simulation creates ecommerce-control, which carries parameter changes. The topics after a run:

ecommerce-control
ecommerce.cdc.dist_centers
ecommerce.cdc.events
ecommerce.cdc.order_items
ecommerce.cdc.orders
ecommerce.cdc.products
ecommerce.cdc.users

Kafka UI listing the six ecommerce.cdc topics and the control topic, with their message counts
The change topics in Kafka UI

Kafka UI&rsquo;s schema registry page listing a key and a value subject for each ecommerce.cdc topic
The Avro subjects in Karapace

In Kafka UI, open a topic such as ecommerce.cdc.orders. Each message is one change, and its op field says which kind:

opMeaning
ra row read in the first snapshot
ca new row (create)
uan update: after holds the row’s new values

Kafka UI showing the newest ecommerce.cdc.orders messages, decoded with their Avro schemas
Change events on the orders topic

This update shows an order moving from Processing to Shipped, with the time it shipped. source records where the change came from: the table, the transaction and its position in the WAL (lsn).

 1{
 2  "before": null,
 3  "after": {
 4    "id": "23b2d4a9-360d-4973-a6da-fd1cc394ca10",
 5    "user_id": "862dd089-c746-4127-89d7-db4fcab59a2e",
 6    "status": "Shipped",
 7    "num_of_items": 1,
 8    "created_at": "2026-09-30T13:38:56+00:00",
 9    "updated_at": "2026-09-30T13:39:01+00:00",
10    "shipped_at": "2026-09-30T13:39:01+00:00",
11    "delivered_at": null,
12    "cancelled_at": null,
13    "returned_at": null
14  },
15  "source": {
16    "version": "3.5.1.Final",
17    "connector": "postgresql",
18    "name": "ecommerce",
19    "ts_ms": 1790775541329,
20    "snapshot": "false",
21    "db": "odctl",
22    "sequence": "[\"161004032\",\"161004832\"]",
23    "ts_us": 1790775541329358,
24    "ts_ns": 1790775541329358000,
25    "schema": "cdc",
26    "table": "orders",
27    "txId": 45922,
28    "lsn": 161004832,
29    "xmin": null,
30    "origin": null,
31    "origin_lsn": null
32  },
33  "transaction": null,
34  "op": "u",
35  "ts_ms": 1790775541830,
36  "ts_us": 1790775541830602,
37  "ts_ns": 1790775541830602065
38}

before is null for updates, because PostgreSQL logs only the new row by default.

In the SeaweedFS file browser, go to odctl-dev/ecommerce-cdc/. There is a folder per topic. Each file holds the change events of one partition, one per line. The first line of the first orders file, 0-0.jsonl, is a create event. The sink adds the message’s key, offset and time around the change event:

 1{
 2  "offset": 0,
 3  "value": {
 4    "before": null,
 5    "after": {
 6      "id": "23b2d4a9-360d-4973-a6da-fd1cc394ca10",
 7      "user_id": "862dd089-c746-4127-89d7-db4fcab59a2e",
 8      "status": "Processing",
 9      "num_of_items": 1,
10      "created_at": "2026-09-30T13:38:56+00:00",
11      "updated_at": "2026-09-30T13:38:56+00:00",
12      "shipped_at": null,
13      "delivered_at": null,
14      "cancelled_at": null,
15      "returned_at": null
16    },
17    "op": "c"
18  },
19  "key": {
20    "id": "23b2d4a9-360d-4973-a6da-fd1cc394ca10"
21  },
22  "timestamp": "2026-09-30T13:38:57.992Z"
23}

The line is shown across several lines here. source, transaction and the ts_* times are left out of value, because they have the same fields as the change event above.

SeaweedFS file browser listing the JSON lines files the sink wrote for the orders topic
Files the S3 sink wrote to SeaweedFS

Step 4: Change the Simulation While It Runs

The simulation reads parameter changes from the Kafka topic ecommerce-control, through dynamic-des’s KafkaIngress. It uses each new value from the next time it draws one. control.py sends a change. With the simulation and both connectors running, cut the pickers from three to one:

1python -m ecommerce.simulation.control ecommerce.resources.pickers.current_cap 1
Sent ecommerce.resources.pickers.current_cap = 1.0 to ecommerce-control

One picker packs about 12 orders a minute, which is fewer than the 20 or so placed, so the queue grows. The Shipped updates in ecommerce.cdc.orders fall from about 20 a minute to about 11. Cancelled updates appear as customers’ patience runs out, which never happens with three pickers. In one run, the change was sent at 12:45:02 UTC. These are the changes per minute in ecommerce.cdc.orders, counted by the minute of each Debezium event:

Minute (UTC)PickersNew ordersShippedCancelled
12:40322220
12:41318180
12:42323220
12:43320210
12:44327270
12:451 from 12:45:0222130
12:46118103
12:47110136
12:48120100
12:49116134
12:50124127
12:5111694
12:521171210
12:53113108

Other parameters work the same way, such as ecommerce.arrival.visitor.rate for visitors a second, or ecommerce.variables.buy for the chance a visitor buys. python -m ecommerce.simulation.control --help lists them all. A change lasts until the simulation stops.

Clean Up

1python -m ecommerce.stores.cleanup

It removes only this project’s objects, and keeps the services running:

  • Connectors: it stops both and deletes their stored offsets, so a new source connector takes a fresh snapshot. It then deletes them.
  • Replication slot: it drops ecommerce_cdc. PostgreSQL keeps the WAL for an unused slot forever, so a slot left behind makes the WAL grow without limit.
  • Tables: it drops the six tables in cdc.
  • Kafka: it deletes the topics starting with ecommerce., and ecommerce-control.
  • SeaweedFS: it deletes the files under odctl-dev/ecommerce-cdc/.

In a test run, the first clean-up deleted both connectors, the seven topics and 45 files. A second run found nothing left to delete, so it is safe to run twice.

To stop the services and delete their data:

1odctl down --all --volumes          # answer y; --volumes also deletes the data
2deactivate
3rm -rf .venv