Estuary

PostgreSQL to Kafka: 3 Ways to Stream Database Changes

Stream PostgreSQL changes to Kafka using managed CDC, Debezium, or JDBC polling. Compare how each method handles inserts, updates, deletes, and Kafka messages.

PostgreSQL to Kafka
Share this article

PostgreSQL data can be sent to Apache Kafka using log-based change data capture (CDC) or query-based polling.

For continuous CDC, tools such as Estuary and Debezium read inserts, updates, and deletes from PostgreSQL’s write-ahead log (WAL) through logical replication and publish those changes to Kafka.

For simpler periodic ingestion, the Kafka Connect JDBC Source Connector can query PostgreSQL on a schedule using timestamp or incrementing columns. This approach doesn't read PostgreSQL WAL and isn't equivalent to log-based CDC.

In this guide, we’ll compare three approaches:

  • Estuary: managed PostgreSQL WAL-based CDC into Kafka
  • Debezium: self-managed PostgreSQL CDC through Kafka Connect
  • JDBC Source Connector: query-based polling for incremental ingestion

PostgreSQL to Kafka: 3 Methods Compared

MethodChange detectionDeletesBest for
EstuaryPostgreSQL WAL / logical replicationYesManaged CDC into Kafka
DebeziumPostgreSQL WAL / logical replicationYesSelf-managed open-source CDC
Kafka Connect JDBC SourceSQL polling using timestamp or incrementing columnsRequires additional handlingPeriodic incremental ingestion

The key difference is how each method detects changes.

Estuary and Debezium read PostgreSQL change events from the WAL, so inserts, updates, and deletes can be captured continuously.

Kafka Connect JDBC Source periodically queries PostgreSQL for rows that are new or have changed since the previous poll. It works well for scheduled ingestion, but deletes are not automatically discoverable once a source row has been removed.

How Does PostgreSQL to Kafka CDC Work?

For log-based CDC, PostgreSQL changes are captured from the write-ahead log (WAL) and published to Kafka as events.

A typical architecture looks like this:

PostgreSQL → WAL / logical replication → CDC connector → Kafka topics → consumers

  1. PostgreSQL records inserts, updates, and deletes in the WAL.
  2. A CDC connector reads those changes through PostgreSQL logical replication.
  3. Each database change is converted into an event.
  4. Events are written to Kafka topics.
  5. Downstream consumers can process those events independently.

PostgreSQL logical replication is the foundation for this pattern. See the PostgreSQL logical replication documentation for how publications, replication slots, and WAL-based change streaming work.

Tools such as Debezium and Estuary’s PostgreSQL connector use PostgreSQL logical replication to capture changes.

The Kafka side then receives those changes as records in topics, where they can be consumed by applications, stream processors, analytics systems, or other downstream services.

This is different from JDBC polling, where the connector periodically runs SQL queries against PostgreSQL and looks for new or updated rows instead of reading the WAL.

Method 1: Stream PostgreSQL Changes to Kafka with Estuary

Estuary can capture PostgreSQL changes from the write-ahead log (WAL) and continuously materialize them into Apache Kafka topics.

The pipeline looks like this:

PostgreSQL → WAL → Estuary capture → Estuary collections → Kafka topics

Estuary’s PostgreSQL source connector uses logical replication to capture database changes. By default, it first backfills the current contents of selected tables and then transitions to ongoing CDC.

Step 1: Prepare PostgreSQL for CDC

The PostgreSQL source requires:

  • wal_level=logical
  • a user with replication permissions
  • a replication slot
  • a publication containing the tables you want to capture
  • network access between Estuary and PostgreSQL

Estuary supports PostgreSQL 10 and later across major cloud platforms and self-hosted deployments.

Step 2: Create the PostgreSQL Capture

PostgreSQL Connector

In Estuary:

  1. Create a PostgreSQL capture.
  2. Enter the source database connection details.
  3. Discover the schemas and tables.
  4. Select the tables you want to stream.
  5. Publish the capture.
Configure Postgres as a Source in Estuary

Unless backfilling is disabled, existing rows are captured first and then new change events continue flowing from PostgreSQL WAL.

Step 3: Configure Kafka as the Destination

Select Kafka Materialization Connector

Use Estuary’s Apache Kafka destination connector to materialize the captured collections into Kafka topics.

You need:

  • Kafka bootstrap servers
  • authentication credentials
  • TLS connectivity
  • a message format
  • a schema registry if you use Avro
Configure Kafka Destination

The connector supports JSON and Avro message formats. JSON does not require a schema registry; Avro does.

Step 4: Map Collections to Kafka Topics

Select the PostgreSQL-derived collections you want to send to Kafka and map them to the appropriate topics.

For each materialization, you can configure settings such as:

  • topic name
  • partition count for newly created topics
  • replication factor
  • field selection
  • JSON or Avro encoding

Once published, changes captured from PostgreSQL continue flowing into the configured Kafka topics.

Delivery Semantics

Estuary’s current Apache Kafka materialization connector uses at-least-once, non-transactional delivery.

That means consumers should be designed to tolerate the possibility of duplicate messages. Estuary does not currently document exactly-once transactional delivery for this connector.

When Is Estuary a Good Fit?

Use Estuary when:

  • you need ongoing PostgreSQL WAL-based CDC into Kafka
  • you want an initial backfill before continuous streaming begins
  • you do not want to operate Debezium and Kafka Connect yourself
  • you want JSON or Avro output
  • you want a managed capture and materialization pipeline

Things to Consider

  • PostgreSQL logical replication and WAL retention still need to be configured and monitored.
  • Kafka delivery is currently at-least-once, so downstream consumers should handle duplicates safely.
  • Avro requires a schema registry.
  • This approach moves PostgreSQL changes into Kafka; your downstream consumers still need to define how those events are interpreted and processed.

Method 2: Stream PostgreSQL Changes to Kafka with Debezium

Postgres Kafka Connection with Debezium
Image Source

Debezium is an open-source CDC platform that can stream PostgreSQL changes into Kafka using PostgreSQL logical replication.

The architecture looks like this:

PostgreSQL → WAL → Debezium PostgreSQL connector → Kafka Connect → Kafka topics

The Debezium PostgreSQL connector reads change events from PostgreSQL’s write-ahead log and publishes them to Kafka through Kafka Connect.

Step 1: Prepare PostgreSQL for Logical Replication

PostgreSQL must be configured for logical replication.

Typical requirements include:

  • wal_level=logical
  • sufficient replication slots
  • sufficient WAL senders
  • a replication user with access to the source tables
  • a publication for the tables you want to capture

Debezium uses PostgreSQL logical decoding to read changes from the replication stream.

Step 2: Run Kafka Connect with the Debezium Connector

Deploy Kafka Connect with the Debezium PostgreSQL connector installed.

Kafka Connect manages the connector process, configuration, offsets, and communication with Kafka.

See the Debezium PostgreSQL connector documentation for the current installation and deployment options rather than relying on version-specific Docker examples embedded in this guide.

Step 3: Configure the PostgreSQL Connector

A Debezium connector configuration typically includes:

  • PostgreSQL hostname and port
  • database name
  • database credentials
  • replication slot
  • publication
  • schemas or tables to capture
  • Kafka topic prefix

Once the connector starts, Debezium can take an initial snapshot of the selected tables and then continue streaming subsequent WAL changes.

Step 4: Consume the Kafka Topics

Debezium writes database change events to Kafka topics.

By default, table changes are typically represented as change-event records that include information about the previous and new row state, operation type, and source metadata.

Downstream consumers can then use those topics for:

  • stream processing
  • cache updates
  • search indexing
  • analytics
  • event-driven applications
  • replication into other systems

When Is Debezium a Good Fit?

Use Debezium when:

  • you want open-source PostgreSQL CDC
  • Kafka is already a core part of your architecture
  • your team is comfortable operating Kafka Connect
  • you need direct control over connector settings and CDC event structure
  • PostgreSQL changes need to feed multiple Kafka consumers

Things to Consider

  • Your team operates Kafka Connect and the Debezium connector.
  • PostgreSQL replication slots and WAL retention need to be monitored.
  • Connector offsets, failures, upgrades, and recovery are your responsibility.
  • Schema changes need to be handled carefully across producers and consumers.
  • The CDC event format is more detailed than a simple row export, so downstream consumers need to understand the Debezium event structure.

Method 3: Send PostgreSQL Data to Kafka with the JDBC Source Connector

The Kafka Connect JDBC Source Connector can move PostgreSQL data into Kafka by running SQL queries on a schedule.

Unlike Estuary and Debezium, it does not read PostgreSQL’s WAL and is not log-based CDC.

The architecture looks like this:

PostgreSQL → SQL polling → JDBC Source Connector → Kafka topics

The connector periodically queries PostgreSQL and publishes the returned rows to Kafka.

How JDBC Polling Works

The JDBC Source Connector supports several ingestion modes:

  • Bulk: reads the full table on each poll
  • Incrementing: reads rows with an increasing numeric column
  • Timestamp: reads rows whose timestamp column has changed
  • Timestamp + incrementing: combines both approaches for more reliable incremental ingestion

See the Kafka Connect JDBC Source Connector documentation for current configuration options.

For example, a timestamp-based setup can track an updated_at column and only query rows changed since the previous poll.

Step 1: Prepare PostgreSQL

Make sure PostgreSQL has:

  • network access from Kafka Connect
  • a database user with read access
  • a reliable incrementing or timestamp column if you want incremental ingestion

For example:

plaintext
updated_at TIMESTAMP

The tracked column should change whenever the source row changes.

Step 2: Configure the JDBC Source Connector

A typical configuration includes:

  • PostgreSQL JDBC connection URL
  • database credentials
  • source table or query
  • ingestion mode
  • tracking column
  • polling interval
  • Kafka topic prefix

For example, a timestamp-based configuration might use:

plaintext
mode=timestamp timestamp.column.name=updated_at topic.prefix=postgres-

The connector then periodically queries PostgreSQL for rows whose tracking value is newer than the last processed value.

Step 3: Publish Rows to Kafka

Each query result is converted into Kafka records and written to the configured topic.

This approach works well when:

  • periodic ingestion is sufficient
  • source tables have reliable tracking columns
  • you already operate Kafka Connect
  • you do not need PostgreSQL WAL-based CDC

What About Deletes?

Deletes are the main limitation of JDBC polling.

Once a row is deleted from PostgreSQL, a normal incremental query cannot return it because the row no longer exists.

If delete propagation matters, you need additional logic such as:

  • soft-delete columns
  • audit tables
  • delete logs
  • triggers
  • or a separate CDC mechanism

This is one of the biggest differences between JDBC polling and WAL-based CDC.

When Is the JDBC Source Connector a Good Fit?

Use it when:

  • updates every few seconds or minutes are acceptable
  • PostgreSQL has reliable timestamp or incrementing columns
  • you already use Kafka Connect
  • you want a simpler polling-based ingestion path

Things to Consider

  • It is polling, not WAL-based CDC.
  • Delete events are not automatically captured.
  • More frequent polling increases query activity on PostgreSQL.
  • Rows can be missed if the tracking column is not updated reliably.
  • Freshness depends on the polling interval.

How Are PostgreSQL Inserts, Updates, and Deletes Handled in Kafka?

How PostgreSQL changes reach Kafka depends on whether the method uses WAL-based CDC or SQL polling.

ChangeEstuaryDebeziumJDBC Source
InsertCaptured from PostgreSQL WAL and written to KafkaCaptured from WAL and emitted as a change eventDetected on the next poll
UpdateCaptured from WAL and materialized to KafkaCaptured from WAL with before/after change dataDetected if the tracking column changes
DeleteCaptured from PostgreSQL WALCaptured from WAL as a delete eventRequires additional handling
Existing rowsInitial backfill before ongoing CDCInitial snapshot before streaming changesRead according to the connector mode

With Estuary and Debezium, deletes can be captured because both methods read PostgreSQL change events from the WAL.

With the JDBC Source Connector, a deleted row no longer exists when the next SQL query runs, so delete propagation usually requires soft deletes, audit tables, or another tracking mechanism.

This is the key distinction: WAL-based CDC captures database events, while JDBC polling discovers rows by querying the current table state.

How Should PostgreSQL Rows Map to Kafka Messages?

For PostgreSQL-to-Kafka pipelines, use a stable PostgreSQL key as the Kafka message key whenever possible.

A typical mapping looks like this:

  • PostgreSQL primary key → Kafka message key
  • Row data or change event → Kafka message value
  • PostgreSQL table → Kafka topic

For example:

plaintext
PostgreSQL table: customers Primary key: customer_id = 123 Kafka topic: customers Kafka message key: 123 Kafka message value: customer record or CDC event

Stable message keys matter because Kafka preserves ordering within a partition for records with the same key. They are also important for compacted topics, where later records with the same key can replace older values.

JSON vs Avro

Kafka records can be serialized in different formats.

JSON is simple to inspect and does not require a schema registry.

Avro provides an explicit schema and is often a better fit when multiple producers and consumers need consistent field definitions. Estuary’s Kafka destination connector supports both JSON and Avro; Avro requires a schema registry.

For Debezium, the message value typically contains a structured CDC event with operation and source metadata in addition to row data. See the Debezium PostgreSQL connector documentation for the current event format.

The main rule is to keep message keys stable and choose a serialization format that downstream consumers can process consistently.

Which PostgreSQL to Kafka Method Should You Use?

Choose based on whether you need managed CDC, self-managed CDC, or periodic polling.

MethodBest forChange detectionMain consideration
EstuaryManaged PostgreSQL CDC into KafkaWAL / logical replicationAt-least-once Kafka delivery; managed capture and materialization
DebeziumSelf-managed open-source CDCWAL / logical replicationYou operate Kafka Connect, connector recovery, and replication infrastructure
JDBC Source ConnectorPeriodic incremental ingestionSQL pollingDeletes require extra handling and freshness depends on the polling interval

Use Estuary when you want PostgreSQL WAL-based CDC into Kafka without operating Debezium and Kafka Connect yourself.

Use Debezium when Kafka is already central to your architecture and your team wants direct control over the CDC stack and event format.

Use the JDBC Source Connector when scheduled polling is sufficient and your tables have reliable timestamp or incrementing columns.

The key architectural choice is WAL-based CDC vs polling. Use CDC when inserts, updates, and deletes need to be captured as database changes occur. Use JDBC polling when periodic ingestion is enough and delete handling can be managed separately.


Need to stream PostgreSQL changes into Kafka continuously?

Start building with Estuary or review the PostgreSQL source connector and Apache Kafka destination connector documentation.

Start streaming your data for free

Build a Pipeline

About the author

Picture of Jeffrey Richman
Jeffrey RichmanData Engineering & Growth Specialist

Jeffrey is a data engineering professional with over 15 years of experience, helping early-stage data companies scale by combining technical expertise with growth-focused strategies. His writing shares practical insights on data systems and efficient scaling.

Streaming Pipelines.
Simple to Deploy.
Simply Priced.
$0.50/GB of data moved + $.14/connector/hour;
50% less than competing ETL/ELT solutions;
<100ms latency on streaming sinks/sources.