CREATE SOURCE: PostgreSQL (Legacy Syntax)

View as Markdown
Disambiguation
This page reflects the legacy syntax, which requires downtime to handle upstream schema changes. For the new syntax which can handle adding or dropping columns to the upstream tables without downtime, see the new reference page.

CREATE SOURCE connects Materialize to an external system you want to read data from, and provides details about how to decode and interpret that data.

Materialize supports PostgreSQL (11+) as a data source. PostgreSQL 16+ is required for connecting Materialize to a physical replica. To connect to a PostgreSQL instance, you first need to create a connection that specifies access and authentication parameters. Once created, a connection is reusable across multiple CREATE SOURCE statements.

WARNING! Before creating a PostgreSQL source, you must set up logical replication in the upstream database. For step-by-step instructions, see the integration guide for your PostgreSQL service: AlloyDB, Amazon RDS, Amazon Aurora, Azure DB, Google Cloud SQL, Self-hosted.
NOTE: Connections using AWS PrivateLink is for Materialize Cloud only.

Syntax

CREATE SOURCE [IF NOT EXISTS] <src_name>
[IN CLUSTER <cluster_name>]
FROM POSTGRES CONNECTION <connection_name> (
  PUBLICATION '<publication_name>'
  [, TEXT COLUMNS ( <col1> [, ...] ) ]
  [, EXCLUDE COLUMNS ( <col1> [, ...] ) ]
)
<FOR ALL TABLES | FOR SCHEMAS ( <schema1> [, ...] ) | FOR TABLES ( <table1> [AS <subsrc_name>] [, ...] )>
[EXPOSE PROGRESS AS <progress_subsource_name>]
[WITH ( <with_option> [, ...] )]
Syntax element Description
<src_name> The name for the source.
IF NOT EXISTS Optional. If specified, do not throw an error if a source with the same name already exists. Instead, issue a notice and skip the source creation.
IN CLUSTER <cluster_name> Optional. The cluster to maintain this source.
CONNECTION <connection_name> The name of the PostgreSQL connection to use in the source. For details on creating connections, check the CREATE CONNECTION documentation page.
PUBLICATION '<publication_name>' The PostgreSQL publication (the replication data set containing the tables to be streamed to Materialize).
TEXT COLUMNS ( <col1> [, …] ) Optional. Decode data as text for specific columns that contain PostgreSQL types that are unsupported in Materialize.
EXCLUDE COLUMNS ( <col1> [, …] ) Optional. Exclude specific columns that cannot be decoded or should not be included in the subsources created in Materialize.
FOR <table_schema_specification>

Specifies which tables to create subsources for. The following <table_schema_specification>s are supported:

Option Description
ALL TABLES Create subsources for all tables in the publication.
SCHEMAS ( <schema1> [, ...] ) Create subsources for specific schemas in the publication.
TABLES ( <table1> [AS <subsrc_name>] [, ...] ) Create subsources for specific tables in the publication.
EXPOSE PROGRESS AS <progress_subsource_name> Optional. The name of the progress collection for the source. If this is not specified, the progress collection will be named <src_name>_progress. For more information, see Monitoring source progress.
WITH (<with_option> [, …])

Optional. The following <with_option>s are supported:

Option Description
RETAIN HISTORY FOR <retention_period> Private preview. This option has known performance or stability issues and is under active development. Duration for which Materialize retains historical data, which is useful to implement durable subscriptions. Accepts positive interval values (e.g. '1hr'). Default: 1s.
TIMESTAMP INTERVAL [=] <interval> The interval at which timestamps are assigned to data read from this source. Accepts positive interval values (e.g. '500ms', '1s'). The value must be between the system parameters min_timestamp_interval and max_timestamp_interval. Default: the value of the default_timestamp_interval system parameter (1s). The interval can also be changed after creation with ALTER SOURCE.

Features

Change data capture

This source uses PostgreSQL’s native replication protocol to continually ingest changes resulting from INSERT, UPDATE and DELETE operations in the upstream database — a process also known as change data capture.

For this reason, you must configure the upstream PostgreSQL database to support logical replication before creating a source in Materialize. For step-by-step instructions, see the integration guide for your PostgreSQL service: AlloyDB, Amazon RDS, Amazon Aurora, Azure DB, Google Cloud SQL, Self-hosted.

Creating a source

To avoid creating multiple replication slots in the upstream PostgreSQL database and minimize the required bandwidth, Materialize ingests the raw replication stream data for some specific set of tables in your publication.

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source')
  FOR ALL TABLES;

When you define a source, Materialize will automatically:

  1. Create a replication slot in the upstream PostgreSQL database (see PostgreSQL replication slots).

    The name of the replication slot created by Materialize is prefixed with materialize_ for easy identification, and can be looked up in mz_internal.mz_postgres_sources.

    SELECT id, replication_slot FROM mz_internal.mz_postgres_sources;
    
       id   |             replication_slot
    --------+----------------------------------------------
     u8     | materialize_7f8a72d0bf2a4b6e9ebc4e61ba769b71
    
  2. Create a subsource for each original table in the publication.

    SHOW SOURCES;
    
             name         |   type
    ----------------------+-----------
     mz_source            | postgres
     mz_source_progress   | progress
     table_1              | subsource
     table_2              | subsource
    

    And perform an initial, snapshot-based sync of the tables in the publication before it starts ingesting change events.

  3. Incrementally update any materialized or indexed views that depend on the source as change events stream in, as a result of INSERT, UPDATE and DELETE operations in the upstream PostgreSQL database.

PostgreSQL replication slots

Each source ingests the raw replication stream data for all tables in the specified publication using a single replication slot. This allows you to minimize the performance impact on the upstream database, as well as reuse the same source across multiple materializations.

💡 Tip:
  • For PostgreSQL 13+, set a reasonable value for max_slot_wal_keep_size to limit the amount of storage used by replication slots.

  • If you stop using Materialize, or if either the Materialize instance or the PostgreSQL instance crash, delete any replication slots. You can query the mz_internal.mz_postgres_sources table to look up the name of the replication slot created for each source.

  • If you delete all objects that depend on a source without also dropping the source, the upstream replication slot remains and will continue to accumulate data so that the source can resume in the future. To avoid unbounded disk space usage, make sure to use DROP SOURCE or manually delete the replication slot.

PostgreSQL schemas

CREATE SOURCE will attempt to create each upstream table in the same schema as the source. This may lead to naming collisions if, for example, you are replicating schema1.table_1 and schema2.table_1. Use the FOR TABLES clause to provide aliases for each upstream table, in such cases, or to specify an alternative destination schema in Materialize.

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source')
  FOR TABLES (schema1.table_1 AS s1_table_1, schema2_table_1 AS s2_table_1);

Reading from a physical standby

Materialize can replicate from a PostgreSQL physical standby (read replica) instead of the primary, using logical decoding on the standby. This requires PostgreSQL 16+ on both the primary and the standby, since earlier versions do not support creating logical replication slots on a standby.

When the upstream is a standby, the replication slot is created on the standby and Materialize only connects to the standby. Note that slot creation on a standby can block until the primary emits a standby snapshot (a RUNNING_XACTS WAL record). On an idle primary, run SELECT pg_log_standby_snapshot() on the primary to unblock source creation.

Monitoring source progress

By default, PostgreSQL sources expose progress metadata as a subsource that you can use to monitor source ingestion progress. The name of the progress subsource can be specified when creating a source using the EXPOSE PROGRESS AS clause; otherwise, it will be named <src_name>_progress.

The following metadata is available for each source as a progress subsource:

Field Type Meaning
lsn uint8 The last Log Sequence Number (LSN) consumed from the upstream PostgreSQL replication stream.

And can be queried using:

SELECT lsn
FROM <src_name>_progress;

The reported LSN should increase as Materialize consumes new WAL records from the upstream PostgreSQL database. For more details on monitoring source ingestion progress and debugging related issues, see Troubleshooting.

Known limitations

Publication membership

PostgreSQL’s logical replication API does not provide a signal when users remove tables from publications. Because of this, Materialize relies on periodic checks to determine if a table has been removed from a publication, at which time it generates an irrevocable error, preventing any values from being read from the table.

However, it is possible to remove a table from a publication and then re-add it before Materialize notices that the table was removed. In this case, Materialize can no longer provide any consistency guarantees about the data we present from the table and, unfortunately, is wholly unaware that this occurred.

To mitigate this issue, if you need to drop and re-add a table to a publication, ensure that you remove the table/subsource from the source before re-adding it using the DROP SOURCE command.

Supported types

Materialize natively supports the following PostgreSQL types (including the array type for each of the types):

  • bool
  • bpchar
  • bytea
  • char
  • date
  • daterange
  • float4
  • float8
  • int2
  • int2vector
  • int4
  • int4range
  • int8
  • int8range
  • interval
  • json
  • jsonb
  • numeric
  • numrange
  • oid
  • text
  • time
  • timestamp
  • timestamptz
  • tsrange
  • tstzrange
  • uuid
  • varchar

Replicating tables that contain unsupported data types is possible via the TEXT COLUMNS option. The specified columns will be treated as text; i.e., will not have the expected PostgreSQL type features. For example:

  • enum: When decoded as text, the implicit ordering of the original PostgreSQL enum type is not preserved; instead, Materialize will sort values as text.

  • money: When decoded as text, resulting text value cannot be cast back to numeric, since PostgreSQL adds typical currency formatting to the output.

Inherited tables

When using PostgreSQL table inheritance, PostgreSQL serves data from SELECTs as if the inheriting tables’ data is also present in the inherited table. However, both PostgreSQL’s logical replication and COPY only present data written to the tables themselves, i.e. the inheriting data is not treated as part of the inherited table.

PostgreSQL sources use logical replication and COPY to ingest table data, so inheriting tables’ data will only be ingested as part of the inheriting table, i.e. in Materialize, the data will not be returned when serving SELECTs from the inherited table.

You can mimic PostgreSQL’s SELECT behavior with inherited tables by creating a materialized view that unions data from the inherited and inheriting tables (using UNION ALL). However, if new tables inherit from the table, data from the inheriting tables will not be available in the view. You will need to add the inheriting tables via ADD SUBSOURCE and create a new view (materialized or non-) that unions the new table.

Handling upstream operations

This section describes how changes to upstream tables that Materialize ingests affect the corresponding Materialize tables.

Adding a column

When you add a new column to your upstream table, Materialize continues to ingest only the existing columns.

To incorporate the new column:

Dropping a column

Dropping columns that Materialize does not ingest (for example, columns added after the source was created, or columns that are excluded) is supported. As these columns were never ingested, you can drop them without issue.

If your Materialize source ingests a column, dropping that column from your upstream table puts the affected table into an error state.

Changing constraints

Materialize ignores the following constraint changes: foreign key, CHECK, and EXCLUSION. As such, you can add or drop them without affecting ingestion.

Materialize also ignores NOT NULL, UNIQUE, and PRIMARY KEY constraints that are added after the Materialize table is created (that is, the table was created without them). Adding such a constraint, and later dropping it, does not affect ingestion.

Dropping a NOT NULL, UNIQUE, or PRIMARY KEY constraint that existed when the table was created puts the affected table into an error state.

If using the new CREATE SOURCE and CREATE TABLE FROM SOURCE syntax, you can safely drop such a constraint by first excluding it in Materialize. See Handle upstream constraint drop.

Changing a column’s data type

Changing an ingested column’s data type upstream puts the affected Materialize table into an error state unless the column was ingested as text via the TEXT COLUMNS option. Ingestion for that table stops, and you must drop and recreate the table in Materialize to resume ingestion.

Renaming a column

Renaming a column that Materialize ingests puts the affected table into an error state. Ingestion for that table stops, and you must drop and recreate the table in Materialize to resume ingestion.

Table-level operations

The following upstream operations put the affected table into an error state. Ingestion for that table stops, and you must drop and recreate the affected table in Materialize to resume:

  • Dropping a table (DROP TABLE), or removing it from the publication (ALTER PUBLICATION ... DROP TABLE).
  • Renaming a table or moving it to a different schema.
  • Setting a table’s replica identity to anything other than FULL (ALTER TABLE ... REPLICA IDENTITY).
  • Truncating a table (TRUNCATE). To clear a table without putting it into an error state, use an unqualified DELETE FROM t; instead.

Source failure states and recovery

Operations that do not require re-creating the source

Materialize tracks a log sequence number (LSN) as it consumes the upstream write-ahead log (WAL), and the source’s replication slot retains the WAL that Materialize has not yet consumed. Because the slot outlives the connection, routine operational events do not lose data: after a transient interruption the source stalls, then resumes from its committed LSN and catches up automatically. No action is required for the following operations:

  • Restarting or patching PostgreSQL (including OS-level restarts).
  • Restarting Materialize. The source resumes from its committed LSN and does not re-snapshot already-ingested data.
  • Transient network interruptions between Materialize and PostgreSQL. These surface as connection closed.
  • Resizing the cluster that hosts the source, or changing its replication factor. Briefly, the source may report replication slot ... is active while the upstream releases the slot from the previous connection.
  • The upstream database running out of disk space, once space is freed.
NOTE: Recovery after an interruption depends on the WAL that the replication slot is holding still being available upstream. An interruption long enough for the slot to be invalidated, or for the slot to be dropped, is not recoverable. See Replication slot invalidated and Replication slot dropped or rewound.
WARNING! While a source is disconnected, the upstream WAL accumulates behind its replication slot and cannot be reclaimed. A long outage, an undersized source cluster, or a source cluster stuck in a restart loop can therefore consume significant upstream disk. Monitor restart_lsn in pg_replication_slots during planned maintenance.

Operations that require re-creating the source

A smaller set of events breaks LSN continuity or destroys the replication slot. When this happens, Materialize cannot guarantee a correct, gap-free view of your data. Most of these put the entire source into an error or permanently stalled state. One, restoring from a volume or disk snapshot, cannot be detected at all, so the source keeps running on diverged data. Every event in this section requires re-creating the source. Upstream changes to an individual table’s schema are handled separately, and do not error the entire source.

In each case below, the remediation is to drop and re-create the source:

DROP SOURCE mz_source CASCADE;

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source');

-- Re-create the tables you were ingesting.
CREATE TABLE table_1 FROM SOURCE mz_source (REFERENCE public.table_1);

If you are using the legacy CREATE SOURCE ... FOR TABLES syntax, re-create the source with the same FOR TABLES list instead of adding tables separately.

Because a re-created source snapshots from the current state of the upstream database, any changes it missed while it was in an error state are reflected in the snapshot rather than replayed as individual updates.

WARNING! CASCADE drops every object that depends on the source, including its tables, views, materialized views, indexes, and sinks. Capture their definitions before you run it, and re-create them once the new source has finished snapshotting.

Point-in-time restore

Restoring the source database from a backup, including restoring to a different server for disaster recovery, increments the PostgreSQL timeline and is detected as a discontinuity. The source fails with an error of the form:

unsupported action: database restored from point-in-time backup. Expected
timeline ID 8 but got 9

The same error covers other events that change the timeline, such as a managed failover between replicas. To see the timeline a source is pinned to, query mz_internal.mz_postgres_sources:

SELECT s.name, p.replication_slot, p.timeline_id
FROM mz_internal.mz_postgres_sources p
JOIN mz_catalog.mz_sources s ON s.id = p.id;

If your upstream fails over between replicas as part of routine maintenance, see High-availability failovers.

Restoring from a volume or disk snapshot

Restoring the upstream data directory from a crash-consistent volume or disk snapshot rolls the database back, but preserves the timeline ID and the replication slot. Materialize cannot detect this kind of restore. The source keeps running without an error, but its contents diverge from upstream. This can surface later as incorrect results, or as negative-accumulation errors in queries such as Non-positive multiplicity.

WARNING! After any restore of this kind, drop and re-create the source even if it reports as running. Do not wait for the source to enter an error state, because it will not.

Promotion of a physical replica

When a source reads from a physical standby (read replica) rather than the primary, promoting that standby to a primary fails the source with:

unsupported action: upstream physical replica status changed (e.g. a physical
replica was promoted to a primary). Expected pg_is_in_recovery()=true but got
false

Materialize detects the promotion while the replication stream is live, without waiting for a restart. Re-create the source against the promoted node.

Replication slot invalidated

PostgreSQL invalidates a replication slot once the WAL it holds exceeds max_slot_wal_keep_size. This protects the upstream from running out of disk, at the cost of ending replication. The source fails with:

replication slot has been invalidated because it exceeded the maximum reserved
size

To avoid this, size the source cluster so that it keeps up with the upstream write rate, and set max_slot_wal_keep_size high enough to cover your longest expected outage. Some hosted PostgreSQL services set this value for you and do not allow it to be raised.

Replication slot dropped or rewound

If the slot Materialize is using is dropped upstream, or the upstream is rebuilt from a base backup (which does not carry replication slots), a new slot starts at the current LSN, past the point the source needs to resume from. The source stalls with:

slot overcompacted. Requested LSN ... but only LSNs >= ... are available

For diagnosis steps, see Slot overcompacted. PostgreSQL refuses to drop a slot that is in use, so this generally happens only while the source is paused or disconnected.

Not every rewind is caught this way. A rewind that leaves the slot able to serve the LSN the source asks for, such as restoring from a volume or disk snapshot, raises no error at all.

Dropping the publication

Running DROP PUBLICATION upstream stalls the source, and all of its tables, with:

publication "mz_source" does not exist

Re-create the publication upstream, then re-create the source.

Major version upgrades

A PostgreSQL major version upgrade rewrites the on-disk format and does not preserve the replication slot, so there is no in-place recovery. To upgrade without a gap in your downstream views, run a second source against the upgraded instance in parallel and cut over once it has hydrated. See Upgrade the major version of your PostgreSQL source.

High-availability failovers

Some managed PostgreSQL services increment the timeline during routine high-availability operations, such as maintenance, a machine-tier change, or an automatic failover between replicas. Materialize cannot distinguish these from a genuine restore, so by default they fail the source with the Expected timeline ID error.

On self-managed Materialize, where the upstream service guarantees that a failover is a contiguous fork of the WAL with no data loss, you can disable timeline validation with the pg_source_validate_timeline system parameter:

ALTER SYSTEM SET pg_source_validate_timeline = false;

This parameter is not available on Materialize Cloud. There, a high-availability failover that changes the timeline requires re-creating the source.

WARNING! Disabling this check is a trade-off. With it off, Materialize also does not detect a genuine point-in-time restore or any other discontinuous timeline change, and silently ingesting across one can corrupt the contents of the source. Only disable it when your provider documents that its failovers preserve WAL continuity for logical replication subscribers, and re-create the source manually after any operation that does not.

Examples

! Important: Before creating a PostgreSQL source, you must set up logical replication in the upstream database. For step-by-step instructions, see the integration guide for your PostgreSQL service: AlloyDB, Amazon RDS, Amazon Aurora, Azure DB, Google Cloud SQL, Self-hosted.

Creating a connection

A connection describes how to connect and authenticate to an external system you want Materialize to read data from.

Once created, a connection is reusable across multiple CREATE SOURCE statements. For more details on creating connections, check the CREATE CONNECTION documentation page.

CREATE SECRET pgpass AS '<POSTGRES_PASSWORD>';

CREATE CONNECTION pg_connection TO POSTGRES (
    HOST 'instance.foo000.us-west-1.rds.amazonaws.com',
    PORT 5432,
    USER 'postgres',
    PASSWORD SECRET pgpass,
    SSL MODE 'require',
    DATABASE 'postgres'
);

If your PostgreSQL server is not exposed to the public internet, you can tunnel the connection through an AWS PrivateLink service (Materialize Cloud) or an SSH bastion host.

CREATE CONNECTION ssh_connection TO SSH TUNNEL (
    HOST 'bastion-host',
    PORT 22,
    USER 'materialize',
);
CREATE CONNECTION pg_connection TO POSTGRES (
    HOST 'instance.foo000.us-west-1.rds.amazonaws.com',
    PORT 5432,
    SSH TUNNEL ssh_connection,
    DATABASE 'postgres'
);

For step-by-step instructions on creating SSH tunnel connections and configuring an SSH bastion server to accept connections from Materialize, check this guide.

Creating a source

Create subsources for all tables included in the PostgreSQL publication

CREATE SOURCE mz_source
    FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source')
    FOR ALL TABLES;

Create subsources for all tables from specific schemas included in the PostgreSQL publication

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source')
  FOR SCHEMAS (public, project);

Create subsources for specific tables included in the PostgreSQL publication

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (PUBLICATION 'mz_source')
  FOR TABLES (table_1, table_2 AS alias_table_2);

Handling unsupported types

If the publication contains tables that use data types unsupported by Materialize, use the TEXT COLUMNS option to decode data as text for the affected columns. This option expects the upstream names of the replicated table and column (i.e. as defined in your PostgreSQL database).

CREATE SOURCE mz_source
  FROM POSTGRES CONNECTION pg_connection (
    PUBLICATION 'mz_source',
    TEXT COLUMNS (upstream_table_name.column_of_unsupported_type)
  ) FOR ALL TABLES;

Handling errors and schema changes

NOTE: Work to more smoothly support ddl changes to upstream tables is currently in progress. The work introduces the ability to re-ingest the same upstream table under a new schema and switch over without downtime.

To handle upstream schema changes or errored subsources, use the DROP SOURCE syntax to drop the affected subsource, and then ALTER SOURCE...ADD SUBSOURCE to add the subsource back to the source.

-- List all subsources in mz_source
SHOW SUBSOURCES ON mz_source;

-- Get rid of an outdated or errored subsource
DROP SOURCE table_1;

-- Start ingesting the table with the updated schema or fix
ALTER SOURCE mz_source ADD SUBSOURCE table_1;

Adding subsources

When adding subsources to a PostgreSQL source, Materialize opens a temporary replication slot to snapshot the new subsources’ current states. After completing the snapshot, the table will be kept up-to-date, like all other tables in the publication.

Dropping subsources

Dropping a subsource prevents Materialize from ingesting any data from it, in addition to dropping any state that Materialize previously had for the table.

Back to top ↑