Skip to main content

Overview

The PostgreSQL CDC streaming source captures real-time changes from PostgreSQL Write-Ahead Log (WAL) using logical replication (pgoutput) and processes them into structured records. It provides comprehensive CDC capabilities with automatic schema evolution tracking and primary key detection.

Real-time CDC

Captures changes as they happen in PostgreSQL

Logical Replication

Uses PostgreSQL’s native logical replication protocol

Primary Key Detection

Automatically identifies and extracts primary key columns

LSN Tracking

Tracks Log Sequence Numbers for reliable resume capabilities

Configuration

Required Parameters

Optional Parameters

SSL Configuration

Record Formats

DB Records Only Format

When return_db_records_only: true (default):
Note: The _mage_* timestamp columns are returned as datetime objects (ISO 8601 format strings when serialized) in UTC timezone, not Unix timestamps. This ensures proper type handling when writing to databases.

Event Types

  • Insert (operation: "INSERT"): INSERT operations
  • Update (operation: "UPDATE"): UPDATE operations with before/after data
  • Delete (operation: "DELETE"): DELETE operations These events contain the actual data changes.
  • Relation (operation: "RELATION"): Table structure metadata
  • Truncate (operation: "TRUNCATE"): TRUNCATE TABLE operations Automatically tracks table schema changes for column name mapping.
  • Begin (operation: "BEGIN"): Transaction begin events
  • Commit (operation: "COMMIT"): Transaction commit events Useful for maintaining transaction boundaries and ensuring data consistency.

Primary Key Detection

The source automatically detects and extracts primary key columns:
  1. Schema Discovery: Queries INFORMATION_SCHEMA for table structure
  2. Primary Key Detection: Identifies actual primary key columns from database constraints
  3. Caching: Stores schema and primary key information for performance
  4. Key Extraction: Extracts primary key column names from row events
Note: key_columns is a list of column names (not values) that can be used for deduplication or upsert operations in downstream sinks.

Examples

Heartbeat Table

The heartbeat table is used to prevent WAL logs from filling up in managed PostgreSQL services (like Aiven). When configured, the source periodically writes to this table to keep the replication slot active.

Schema

The heartbeat table is automatically created if it doesn’t exist. It has the following schema:
Columns:
  • id (INTEGER PRIMARY KEY): Primary key column with default value of 1
  • timestamp (TIMESTAMP): Timestamp column that gets updated with the current timestamp on each heartbeat write

How It Works

  1. Automatic Creation: If the table doesn’t exist, it’s automatically created when the source initializes
  2. Periodic Updates: The source updates the timestamp column every heartbeat_interval_seconds (default: 60 seconds)
  3. Default UPSERT Pattern: Uses INSERT ... ON CONFLICT DO UPDATE to update the single row with id=1, or falls back to UPDATE/INSERT if the database doesn’t support ON CONFLICT
  4. Custom Query: You can override the default heartbeat update query by providing heartbeat_update_query in the configuration

Security Warning for Custom Heartbeat Query

Important Security Notice: The heartbeat_update_query parameter executes SQL queries directly without sanitization or parameterization. This could introduce SQL injection vulnerabilities if the query is sourced from untrusted input.Best Practices:
  • Only use queries from trusted, validated sources
  • Never construct queries dynamically from user input or external APIs
  • Validate that the query is not empty or whitespace-only before use
  • Consider using the default heartbeat mechanism unless you have specific requirements
  • Review and test custom queries thoroughly before deploying to production
The source performs basic validation to ensure the query is not empty or whitespace-only, but does not perform SQL injection protection. Always ensure your configuration management system properly secures and validates this parameter.

Manual Creation

If you prefer to create the table manually, you can use:
Note: The table must have a primary key or unique constraint on the id column for the UPSERT pattern to work correctly.

Prerequisites

PostgreSQL Server Configuration

Create Replication Slot

Create Publication

User Permissions

Troubleshooting

  1. Replication Slot Not Found: Create the replication slot using pg_create_logical_replication_slot()
  2. Publication Not Found: Create the publication using CREATE PUBLICATION
  3. Permission Denied: Ensure user has REPLICATION privilege and SELECT on tables
  4. WAL Level Not Logical: Set wal_level = logical and restart PostgreSQL
  5. Connection Timeout: Increase connect_timeout value
  6. LSN Format Error: Ensure LSN is in format ‘X/Y’ (e.g., ‘0/1234567’)
Check replication slot status:
Check publication status:
Monitor key metrics:
  • Batch Size: Average events per batch
  • Flush Rate: How often batches are flushed
  • Error Rate: Failed events or connections
  • Lag: Time between event occurrence and processing
  • LSN Progress: Track Log Sequence Number progress
Check replication lag:

Integration with Generic IO Sink

The PostgreSQL CDC source works seamlessly with the Generic IO sink, which provides:
  • Automatic Column Type Mapping: Automatically maps _mage_* timestamp columns to appropriate database types (TIMESTAMP, DATETIME2, DateTime64, etc.) based on the target database
  • Metadata Interpolation: Use metadata values from PostgreSQL CDC events in sink configurations for dynamic routing and table naming

Supported Databases

Generic IO Sink supports the following databases:
  • BigQuery
  • ClickHouse
  • DuckDB
  • MySQL
  • MSSQL
  • Postgres

Metadata Interpolation

You can use metadata values from PostgreSQL CDC events in your sink configuration using Python string formatting syntax:
  • {schema}: Schema name from the event
  • {table}: Table name from the event
  • {key_columns}: List of primary key column names (e.g., ["id"] or ["user_id", "tenant_id"])
When using {key_columns} in unique_constraints, it will be automatically converted from a string representation to a list. The format supports both Python-style ("['id']") and JSON-style ('["id"]') array strings.

Example Configurations

How Metadata Interpolation Works

  1. Message Grouping: Messages are automatically grouped by their interpolated config values. For example, if you use table_name: "{schema}_{table}", messages from public.users will be grouped together and written to public_users table.
  2. Key Columns Interpolation: When using {key_columns} in unique_constraints, the sink automatically:
    • Converts the list to a string representation during interpolation
    • Parses it back to a list (supports both "['id']" and '["id"]' formats)
    • Uses it for upsert operations based on unique_conflict_method
Example: If you have a table users with primary key id, and you configure unique_constraints: "{key_columns}", it will automatically use ["id"] for upsert operations.

Best Practices

  1. Use Replication Slots: Always use replication slots for reliable CDC tracking
  2. Monitor WAL Growth: Monitor WAL size and ensure proper cleanup
  3. Filter Events: Use schema/table filters to reduce processing overhead
  4. Monitor Resources: Watch memory and CPU usage
  5. Test Resume: Verify checkpoint functionality with LSN tracking
  6. Secure Connections: Use SSL in production
  7. Regular Backups: Backup checkpoint files
  8. Schema Validation: Test with schema changes
  9. Use Generic IO Sink: Leverage automatic type mapping and metadata interpolation for flexible data routing
  10. Metadata Interpolation: Use {schema}, {table}, and {key_columns} for dynamic table routing and upsert configuration
  11. Heartbeat Tables: Use heartbeat tables for managed services (like Aiven) to prevent WAL log filling

Timestamp Handling

The PostgreSQL CDC source automatically adds Mage timestamp columns:
  • _mage_created_at: Set to event timestamp (datetime) for INSERT operations
  • _mage_updated_at: Set to event timestamp (datetime) for UPDATE operations
  • _mage_deleted_at: Set to event timestamp (datetime) for DELETE operations
All timestamps are in UTC timezone and returned as datetime objects, which ensures:
  • Proper type handling in downstream databases
  • Automatic type conversion in Generic IO sink
  • Consistent timezone handling across systems

Limitations

  • PostgreSQL 10+: Requires PostgreSQL 10 or later for logical replication
  • WAL Level: Requires wal_level = logical
  • Network Dependency: Requires stable network connection
  • Memory Usage: Schema caching uses memory
  • WAL Retention: Depends on PostgreSQL WAL retention settings
  • Timezone: All timestamps are in UTC timezone
  • Replication Slot: Requires a dedicated replication slot per consumer
  • Publication: Tables must be added to a publication to be replicated