Realtime synchronization from PostgreSQL
This topic shows how to stream data from PostgreSQL into StarRocks with a Flink CDC pipeline. The pipeline copies the existing rows of each table, then keeps applying every INSERT, UPDATE, and DELETE from PostgreSQL within seconds. It creates the StarRocks tables for you, so there is no separate schema migration step.
You can load data into StarRocks tables only as a user who has the INSERT privilege on those StarRocks tables. If you do not have the INSERT privilege, follow the instructions provided in GRANT to grant the INSERT privilege to the user that you use to connect to your StarRocks cluster. The syntax is GRANT INSERT ON TABLE <table_name> IN DATABASE <database_name> TO { ROLE <role_name> | USER <user_identity>}.
How it works​
A Flink CDC pipeline is a YAML file with a source (PostgreSQL), a sink (StarRocks), and optional route and transform rules. You submit the file to a Flink cluster, and it runs as a Flink job:
- The PostgreSQL source takes a snapshot of each matching table, then reads changes from the PostgreSQL write-ahead log (WAL) through a logical replication slot.
- The StarRocks sink creates a Primary Key table for each source table that does not already exist in StarRocks, then loads the rows with Stream Load. Because every table has a primary key, updates overwrite the existing row and deletes remove it.
Delivery is at-least-once: after a failure, some changes can be written again, and the Primary Key table makes the repeat harmless.
The pipeline does not synchronize schema changes from PostgreSQL. See Change a table's schema.
Before you begin​
You need:
- PostgreSQL 10 or later. This topic uses the built-in
pgoutputlogical decoding plugin, so no server extension is required. - A StarRocks cluster.
- A Flink 1.20 cluster, and Java 11 or 17 to run it.
- Network access from every Flink TaskManager to PostgreSQL, on its port (default
5432), and to StarRocks, as described in Connect to StarRocks. - A primary key on every PostgreSQL table you synchronize. StarRocks creates the destination table with the same primary key.
Step 1: Prepare PostgreSQL​
Run these steps as a PostgreSQL superuser, or as the owner of the tables. The examples synchronize the tables in the public schema of a database named shop.
-
Enable logical replication. Set
wal_leveltologicalinpostgresql.conf:wal_level = logical
# Each running pipeline uses one replication slot and one WAL sender.
max_replication_slots = 4
max_wal_senders = 4Changing
wal_levelrequires a restart of PostgreSQL. For a managed service such as Amazon RDS or Cloud SQL, set the equivalent parameter (for example,rds.logical_replication = 1) in the instance's parameter group instead.Check the setting after the restart:
SHOW wal_level;wal_level
-----------
logical -
Create a user for Flink CDC. It needs the
REPLICATIONattribute and read access to the tables:CREATE ROLE flink_cdc WITH LOGIN REPLICATION PASSWORD '<password>';
GRANT CONNECT ON DATABASE shop TO flink_cdc;
GRANT USAGE ON SCHEMA public TO flink_cdc;
GRANT SELECT ON ALL TABLES IN SCHEMA public TO flink_cdc;If
pg_hba.confrestricts replication connections, allow this user to connect from the Flink TaskManager hosts. For example:# TYPE DATABASE USER ADDRESS METHOD
host shop flink_cdc 10.0.0.0/24 scram-sha-256
host replication flink_cdc 10.0.0.0/24 scram-sha-256 -
Create a publication for the tables to synchronize. It must list every table that the pipeline reads (the
tablesoption in Step 4):CREATE PUBLICATION flink_pub FOR TABLE public.orders, public.customers;tipCreate the publication yourself. If none exists, the connector tries to run
CREATE PUBLICATION ... FOR ALL TABLES, which only a superuser can do. For any other user that fails, and Flink retries until the error it reports is aboutmax_wal_sendersrather than the missing publication. -
Set the replica identity of each table to
FULL, so that the WAL records the complete previous row for every UPDATE and DELETE:ALTER TABLE public.orders REPLICA IDENTITY FULL;
ALTER TABLE public.customers REPLICA IDENTITY FULL;Do this before you start the pipeline. With the default replica identity, the first UPDATE or DELETE fails the Flink job with a
NullPointerException, and the job keeps failing on that change after every restart.
Step 2: Prepare StarRocks​
Create the destination database and a user for the pipeline:
CREATE DATABASE shop;
CREATE USER flink_sink IDENTIFIED BY '<password>';
GRANT CREATE TABLE ON DATABASE shop TO USER flink_sink;
GRANT SELECT, INSERT, UPDATE, DELETE, ALTER ON ALL TABLES IN DATABASE shop TO USER flink_sink;
Step 3: Install Flink CDC​
-
Download and start Flink 1.20. See First steps in the Flink documentation.
-
Turn on checkpointing by adding this line to the Flink configuration file,
conf/config.yaml, and then restart the cluster:execution.checkpointing.interval: 10sThe PostgreSQL source reports its position in the WAL back to PostgreSQL as checkpoints complete. Without checkpoints, the replication slot keeps every WAL segment since the pipeline started, and the PostgreSQL disk fills up.
-
Download the Flink CDC release that matches your Flink version from the Flink CDC downloads page, and extract it. Starting with Flink CDC 3.6, each release is built for one Flink version, so for Flink 1.20 download
flink-cdc-3.6.0-1.20-bin.tar.gz, not the-2.2package.tar -xzf flink-cdc-3.6.0-1.20-bin.tar.gz
cd flink-cdc-3.6.0-1.20 -
Download the PostgreSQL and StarRocks pipeline connectors from Maven Central into the
libdirectory of Flink CDC. Use the same version as the Flink CDC release:cd lib
wget https://repo1.maven.org/maven2/org/apache/flink/flink-cdc-pipeline-connector-postgres/3.6.0-1.20/flink-cdc-pipeline-connector-postgres-3.6.0-1.20.jar
wget https://repo1.maven.org/maven2/org/apache/flink/flink-cdc-pipeline-connector-starrocks/3.6.0-1.20/flink-cdc-pipeline-connector-starrocks-3.6.0-1.20.jar
cd ..
Step 4: Define the pipeline​
Create a file named postgres-to-starrocks.yaml:
source:
type: postgres
name: PostgreSQL source
hostname: <postgres_host>
port: 5432
username: flink_cdc
password: <password>
# Comma-separated <database>.<schema>.<table> entries; the table part can be a
# regular expression. List the same tables as the publication.
tables: shop.public.orders, shop.public.customers
# Replication slot to create and use. Lowercase letters, digits, and underscores only.
slot.name: flink_starrocks
decoding.plugin.name: pgoutput
debezium.publication.name: flink_pub
sink:
type: starrocks
name: StarRocks sink
# For the StarRocks addresses, see "Connect to StarRocks" below.
jdbc-url: jdbc:mysql://<fe_host>:9030
load-url: <fe_host>:8030
username: flink_sink
password: <password>
# The default is 300000 (5 minutes). A lower value makes changes visible sooner.
sink.buffer-flush.interval-ms: 5000
route:
# Write each table in the public schema to the StarRocks database shop.
- source-table: public.\.*
sink-table: shop.<>
replace-symbol: <>
pipeline:
name: PostgreSQL to StarRocks
parallelism: 1
Pay attention to these settings:
tablesmust match the publication. The snapshot copies every table thattablesselects, but PostgreSQL streams changes only for the tables in the publication. A table outside the publication is copied once and then never updated, and the job reports no error.tablesnames tables as<database>.<schema>.<table>. Inside the pipeline, however, a PostgreSQL table is identified only by<schema>.<table>, sorouteandtransformrules matchpublic.orders, notshop.public.orders.- Without the
routeblock, the sink writes each table to a StarRocks database named after the PostgreSQL schema, so the tables inpublicland in a database calledpublic. sink.buffer-flush.interval-msdefaults to 5 minutes. With the default, a new pipeline looks idle for several minutes after it starts.- To set properties on the tables that the sink creates, add
table.create.properties.<property>to thesinksection, for exampletable.create.properties.replication_num: 1.
For every source and sink option, see the Flink CDC Postgres connector and StarRocks connector references.
Connect to StarRocks​
Set the StarRocks connection options to addresses of the FE:
jdbc-url:jdbc:mysql://<fe_host>:<fe_query_port>. The query port defaults to9030.load-url:<fe_host>:<fe_http_port>. The HTTP port defaults to8030.
For <fe_host>, use an address that stays the same when an FE restarts or is replaced, such as a DNS name or a load balancer in front of the FEs, not the IP address of an FE. If the address stops leading to an FE, the Flink job can no longer load data.
The Flink TaskManagers must be able to reach these FE ports, and also the HTTP port (default 8040) of every BE or CN, because the FE redirects each load request to a BE or CN.
Columns with a function as the default value​
The sink copies each column's default value from PostgreSQL into the StarRocks CREATE TABLE statement. A default that is a function call, such as DEFAULT now(), is not valid in StarRocks, and the job fails with an error like this one:
Invalid default value for 'created_at': date literal [now()] is invalid.
Creating the StarRocks table yourself beforehand does not avoid the error, because the sink still sends its CREATE TABLE statement. Instead, add a transform rule that recomputes the column. A computed column has no default value:
transform:
- source-table: public.orders
projection: order_id, customer, amount, status, CAST(created_at AS TIMESTAMP(6)) AS created_at
List every column of the table in projection. A column that is not listed is not synchronized.
Step 5: Start the pipeline​
Set FLINK_HOME to your Flink installation, and submit the pipeline from the Flink CDC directory:
export FLINK_HOME=/path/to/flink-1.20
bin/flink-cdc.sh postgres-to-starrocks.yaml
Pipeline has been submitted to cluster.
Job ID: cda97bac3b504442d8106a2d70c6074c
Job Description: PostgreSQL to StarRocks
The job appears in the Flink web UI (by default http://<jobmanager_host>:8081). Its status should stay RUNNING. If it changes to RESTARTING, open the job's Exceptions tab, and see Troubleshooting.
Step 6: Verify the synchronization​
-
After the snapshot finishes, the existing rows are in StarRocks:
SELECT * FROM shop.orders ORDER BY order_id; -
Change some data in PostgreSQL:
INSERT INTO public.orders (order_id, customer, amount, status) VALUES (5, 'erin', 42.00, 'new');
UPDATE public.orders SET status = 'returned' WHERE order_id = 3;
DELETE FROM public.orders WHERE order_id = 4; -
Within a few seconds (the flush interval plus the load time), run the query in StarRocks again. It returns the same rows as the same query in PostgreSQL.
Manage the pipeline​
Monitor the replication slot​
A replication slot keeps WAL on the PostgreSQL server until the pipeline confirms that it has read it. While the pipeline is stopped, or if it falls behind, WAL builds up. Check how much WAL the slot is holding:
SELECT slot_name,
active,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)) AS retained_wal
FROM pg_replication_slots;
If you stop a pipeline for good, drop its slot so that PostgreSQL can release the WAL:
SELECT pg_drop_replication_slot('flink_starrocks');
Change a table's schema​
The pipeline does not pick up schema changes from PostgreSQL. If you add a column in PostgreSQL, the job keeps running, but it does not write the new column, even after you add the column to the StarRocks table.
To synchronize a table after a schema change:
- Make the same change to the StarRocks table, for example with ALTER TABLE.
- Take a new snapshot.
Take a new snapshot​
To copy the tables again from the current state of PostgreSQL and then resume streaming:
-
Cancel the Flink job.
-
Drop the replication slot, as shown in Monitor the replication slot.
-
Truncate every StarRocks table that the pipeline writes to:
TRUNCATE TABLE shop.orders;
TRUNCATE TABLE shop.customers; -
Submit the pipeline again. The pipeline takes a new snapshot of every table and then continues streaming changes.
Do not skip the truncate. Changes made in PostgreSQL between cancelling the job and the new snapshot are not replayed. The snapshot overwrites the rows that still exist, but a row deleted in PostgreSQL during that time would stay in StarRocks. The truncated tables are empty, or only partly loaded, until the snapshot finishes.
Troubleshooting​
When the job is RESTARTING, the Exceptions tab of the Flink web UI shows the cause. Look for the innermost Caused by line.
| Error | Cause and solution |
|---|---|
number of requested standby connections exceeds max_wal_senders | Often a symptom of an earlier failure, because every retry opens a replication connection. Check the PostgreSQL server log for the first error. A common one is permission denied for database on CREATE PUBLICATION, which means the publication does not exist. See Step 1. |
NullPointerException in DebeziumSchemaDataTypeInference or extractBeforeDataRecord | A table does not have REPLICA IDENTITY FULL. Set it, and then take a new snapshot. |
Invalid default value for '<column>' | A column's PostgreSQL default is a function. See Columns with a function as the default value. |
Connect to <host>:8040 ... Connection refused | The TaskManager cannot reach the BE or CN that the FE redirected the load to. Make sure every BE and CN HTTP port is reachable from the TaskManagers at the address the BE or CN is registered with. |
Nothing arrives in StarRocks, but the job is RUNNING | Wait for the flush interval. The default sink.buffer-flush.interval-ms is 5 minutes. |
Rows arrive in a database named public | The pipeline has no route rule. See Step 4. |