The SelectDB connector integrates Realtime Compute for Apache Flink with ApsaraDB for SelectDB, a fully managed, Apache Doris-compatible real-time data warehouse on Alibaba Cloud. Use it to build real-time pipelines that read from, write to, or look up data in SelectDB, and to run full database synchronization in YAML-based data ingestion jobs.
Supported capabilities:
| Category | Details |
|---|---|
| Table types | Source table, sink table, dimension table, data ingestion sink |
| Running mode | Stream and batch |
| Data format | JSON and CSV |
| API type | DataStream, SQL, and data ingestion YAML jobs |
| Update/Delete support | Yes |
| Monitoring metrics | None |
Key features:
-
Full database data synchronization
-
Exactly-once semantics via two-phase commit (2PC) — no duplicate or lost records
-
Compatible with Apache Doris 1.0 and later
Prerequisites
Before you begin, make sure you have:
-
Realtime Compute for Apache Flink with Ververica Runtime (VVR) 8.0.10 or later
-
An ApsaraDB for SelectDB instance. See Create an instance.
-
An IP address whitelist configured on the instance. See Configure a whitelist.
Set up the connector
The SelectDB connector is built into VVR 11.1 and later — no manual installation required.
For VVR 8.0.10 through 11.0, install the connector manually:
-
Download the JAR package from Maven Central (Flink versions 1.15–1.17).
-
Upload the JAR to your Realtime Compute for Apache Flink development console. See Manage custom connectors.
-
Reference the connector in your SQL job using
'connector' = 'doris'.
SQL
Syntax
All three table types — source, sink, and dimension — share the same DDL syntax. Specify the table role through the parameters you include.
To use SelectDB as a source table, enable direct cluster connection first. In the ApsaraDB for SelectDB console, go to Instance Details > Network Information and click Enable Direct Cluster Connection. This activates the Arrow Flight SQL protocol for high-throughput parallel reads.
CREATE TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
Parameters
General
| Parameter | Required | Default | Description |
|---|---|---|---|
connector |
Yes | — | Fixed to doris. |
fenodes |
Yes | — | HTTP endpoint of the SelectDB instance: <VPC Address or Public Address>:<HTTP Protocol Port>. Get both from Instance Details > Network Information in the SelectDB console. Example: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080. |
jdbc-url |
No | — | Java Database Connectivity (JDBC) connection string for dimension table lookups and metadata queries: jdbc:mysql://<VPC Address or Public Address>:<MySQL Protocol Port>. Example: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030. |
table.identifier |
Yes | — | Target table in <database>.<table> format. Example: db.tbl. |
username |
Yes | — | Database username. Reset the password from the upper-right corner of the Instance Details page if needed. |
password |
Yes | — | Password for the database username. |
doris.request.retries |
No | 3 |
Number of retries for failed requests. |
doris.request.connect.timeout |
No | 30s |
Connection timeout. |
doris.request.read.timeout |
No | 30s |
Read timeout. |
Source table
| Parameter | Required | Default | Description |
|---|---|---|---|
doris.request.query.timeout |
No | 21600s |
Query timeout (6 hours by default). |
doris.request.tablet.size |
No | 1 |
Number of tablets per partition. Lower values increase Flink parallelism but add more pressure on the database. |
doris.batch.size |
No | 4064 |
Maximum rows read from a Backend (BE) node per request. Increase to reduce connection overhead and network latency. |
doris.exec.mem.limit |
No | 8192mb |
Memory limit per query in bytes (8 GB by default). |
source.use-flight-sql |
No | false |
No configuration needed — enabling Direct Cluster Connection in the SelectDB console activates Arrow Flight SQL automatically. |
source.flight-sql-port |
No | — | Arrow Flight SQL port (arrow_flight_sql_port) of the Frontend (FE) node. |
Sink table
Write mode affects delivery guarantees and flush behavior. Choose based on your consistency requirements:
| Streaming write | Batch write | |
|---|---|---|
| Trigger condition | Follows Flink checkpoint intervals | Periodic flush by data volume or time threshold |
| Delivery guarantee | Exactly-once (via 2PC) | At-least-once; achieve idempotence with the Unique model |
| Latency | Bounded by checkpoint interval | Flexible, independent of checkpoints |
| Fault tolerance | Full Flink state recovery | Relies on Unique model deduplication |
| Parameter | Required | Default | Description |
|---|---|---|---|
sink.label-prefix |
No | — | Label prefix for Stream Load imports. Must be globally unique across all jobs — the same label can only be committed once. Required to guarantee exactly-once semantics across job restarts. |
sink.properties.* |
No | — | Stream Load import parameters passed directly to the SelectDB Stream Load API. See examples below. |
sink.enable-delete |
No | true |
Propagate DELETE operations. Requires the Doris table to have batch deletion enabled and only works with the Unique model. |
sink.enable-2pc |
No | true |
Enable two-phase commit (2PC) for exactly-once semantics. See Explicit Transaction Operations. |
sink.buffer-size |
No | 1 MB |
Write cache buffer size in bytes. Leave at the default. |
sink.buffer-count |
No | 3 |
Number of write cache buffers. Leave at the default. |
sink.max-retries |
No | 3 |
Maximum retries after a commit failure. |
sink.enable.batch-mode |
No | false |
Switch to batch write mode. Flush is controlled by the three sink.buffer-flush.* parameters below instead of checkpoints. Exactly-once is not guaranteed; use the Unique model for idempotence. |
sink.flush.queue-size |
No | 2 |
Cache queue size in batch mode. |
sink.buffer-flush.max-rows |
No | 500000 |
Maximum rows per flush in batch mode. |
sink.buffer-flush.max-bytes |
No | 100 MB |
Maximum bytes per flush in batch mode. |
sink.buffer-flush.interval |
No | 10s |
Flush interval in batch mode. |
sink.ignore.update-before |
No | true |
Ignore update-before events from Flink CDC. |
**sink.properties.* examples:**
CSV format:
'sink.properties.column_separator' = ','
-- If values may contain commas, use a non-printable separator:
-- 'sink.properties.column_separator' = '\x01'
JSON format:
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true'
-- Alternatively: 'sink.properties.strip_outer_array' = 'true'
Dimension table
| Parameter | Required | Default | Description |
|---|---|---|---|
lookup.cache.max-rows |
No | -1 |
Maximum rows in the lookup cache. -1 disables caching. |
lookup.cache.ttl |
No | 10s |
Cache entry time-to-live (TTL). |
lookup.max-retries |
No | 1 |
Retries after a lookup query fails. |
lookup.jdbc.async |
No | false |
Enable asynchronous lookup. |
lookup.jdbc.read.batch.size |
No | 128 |
Maximum batch size per query in async lookup mode. |
lookup.jdbc.read.batch.queue-size |
No | 256 |
Intermediate buffer queue size in async lookup mode. |
lookup.jdbc.read.thread-size |
No | 3 |
JDBC lookup threads per task in async lookup mode. |
Examples
Source table
CREATE TEMPORARY TABLE selectdb_source (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****'
);
Sink table
CREATE TEMPORARY TABLE selectdb_sink (
order_id BIGINT,
user_id BIGINT,
total_amount DECIMAL(10, 2),
order_status TINYINT,
create_time TIMESTAMP(3),
product_name STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'table.identifier' = 'shop_db.orders',
'username' = 'admin',
'password' = '****',
'sink.label-prefix' = 'flink_orders' -- Must be globally unique across jobs
);
Dimension table
SelectDB acts as a lookup dimension table joined against a streaming fact table.
-- Fact table from Kafka
CREATE TEMPORARY TABLE fact_table (
`id` BIGINT,
`name` STRING,
`city` STRING,
`process_time` AS proctime()
) WITH (
'connector' = 'kafka',
...
);
-- Dimension table from SelectDB
CREATE TEMPORARY TABLE dim_city (
`city` STRING,
`level` INT,
`province` STRING,
`country` STRING
) WITH (
'connector' = 'doris',
'fenodes' = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
'jdbc-url' = 'jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030',
'table.identifier' = 'dim.dim_city',
'username' = 'admin',
'password' = '****'
);
-- Temporal join
SELECT a.id, a.name, a.city, c.province, c.country, c.level
FROM fact_table a
LEFT JOIN dim_city FOR SYSTEM_TIME AS OF a.process_time AS c
ON a.city = c.city;
Data ingestion
Use the SelectDB connector as a sink in YAML-based data ingestion jobs for full database synchronization.
Syntax
source:
type: <source-type>
sink:
type: doris
name: Doris Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
username: root
password: ""
Parameters
| Parameter | Required | Default | Description |
|---|---|---|---|
type |
Yes | — | Fixed to doris. |
name |
No | — | Descriptive name for the sink. |
fenodes |
Yes | — | HTTP endpoint: <VPC Address or Public Address>:<HTTP Protocol Port>. Get both from Instance Details > Network Information in the SelectDB console. Example: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080. |
jdbc-url |
No | — | JDBC connection string. Example: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030. |
username |
Yes | — | Database username. |
password |
Yes | — | Password for the database username. |
sink.enable.batch-mode |
No | true |
Batch mode is on by default in data ingestion jobs. Flush is controlled by the three sink.buffer-flush.* parameters. Exactly-once is not guaranteed; use the Unique model for idempotence. |
sink.flush.queue-size |
No | 2 |
Cache queue size. |
sink.buffer-flush.max-rows |
No | 500000 |
Maximum rows per flush. |
sink.buffer-flush.max-bytes |
No | 100 MB |
Maximum bytes per flush. |
sink.buffer-flush.interval |
No | 10s |
Flush interval. Minimum: 1s. |
sink.properties.* |
No | — | Stream Load import parameters. |
**sink.properties.* examples:**
CSV format:
sink.properties.column_separator: ','
# If values may contain commas, use a non-printable separator:
# sink.properties.column_separator: '\x01'
JSON format:
sink.properties.format: 'json'
sink.properties.read_json_by_line: 'true'
Examples
The following examples cover common data ingestion scenarios. Replace the placeholder endpoints, usernames, and passwords with your actual values.
${secret_values.variable_name} references a secret variable that is created in your workspace. For more information, see Manage variables. For a full description of the source parameters, see MySQL YAML connector and Postgres CDC connector.
The connector does not create destination databases. Before you start, create the destination databases that are used in the following examples in your ApsaraDB for SelectDB instance. Destination tables are created automatically if they do not exist.
Single-table synchronization
Synchronize a single MySQL table to SelectDB. The destination table is created automatically if it does not exist.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
# (Optional) Synchronize the existing full data first, and then the incremental data.
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: MySQL to SelectDB Pipeline
Full database synchronization
Synchronize all tables in a MySQL database to SelectDB. The downstream database and table names are the same as the upstream ones. Destination tables are created automatically if they do not exist.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
# Matches all tables in the test_db database.
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
# (Optional) Synchronize tables that are created during the incremental phase without restarting the job.
scan.binlog.newly-added-table.enabled: true
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
# (Optional) The flush interval for batch writing. Lower this value for small workloads so that data is not held in the buffer for too long.
sink.buffer-flush.interval: 10s
pipeline:
name: MySQL to SelectDB Pipeline
Excluding tables from full database synchronization
Skip the tables that you do not want to synchronize, such as temporary tables and log tables. Without a route module, data is written to the destination table that has the same name as the source table.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.\.*
# Tables that match this regular expression are not synchronized. Separate multiple regular expressions with commas (,).
tables.exclude: test_db.tmp_\.*, test_db.log_\.*
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: MySQL to SelectDB Pipeline
Synchronizing to specified databases and tables
Use the route module to rename destinations when the downstream database or table names differ from the upstream ones. The following example synchronizes the tables in test_db to the ods_db database and adds the ods_ prefix to each table name. For example, test_db.orders is synchronized to ods_db.ods_orders.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.\.*
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
route:
# <> is a placeholder that is replaced with the matched source table name.
- source-table: test_db.\.*
sink-table: ods_db.ods_<>
replace-symbol: <>
pipeline:
name: MySQL to SelectDB Pipeline
Merging sharded tables
Merge multiple shard tables that share the same schema into a single SelectDB table. The following example uses the transform module to append a column that identifies the source database and table, and then uses this column with the business primary key as a composite primary key so that you can trace the origin of the merged data.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
# Matches all shards, such as order_db_1.orders_1 and order_db_2.orders_2.
tables: order_db_\d+.orders_\d+
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
# Set this parameter to false when the downstream primary key differs from the upstream one. Otherwise, the rows that correspond to the old primary key are not deleted.
sink.ignore.update-before: false
transform:
- source-table: order_db_\d+.orders_\d+
# Append the src_table column that identifies the source database and table.
projection: "*, __schema_name__ || '.' || __table_name__ AS src_table"
# The composite primary key must be the same as the key columns of the destination table.
primary-keys: order_id, src_table
description: Append the source identifier and set a composite primary key
route:
# Merge all shards into a single table.
- source-table: order_db_\d+.orders_\d+
sink-table: dw_db.merged_orders
pipeline:
name: MySQL sharding to SelectDB Pipeline
Filtering data and pruning columns
Use the transform module to filter data and prune columns. The filter condition is evaluated on the source table fields to decide whether a record is synchronized, and the projection decides which columns are written to the downstream table.
source:
type: mysql
name: MySQL Source
hostname: <yourHostname>
port: 3306
username: ${secret_values.mysql_username}
password: ${secret_values.mysql_password}
tables: test_db.test_source_table
server-id: 5401-5499
scan.startup.mode: initial
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
transform:
- source-table: test_db.test_source_table
# Synchronize only the records whose order_status is PAID.
filter: "order_status = 'PAID'"
# Write only the following four columns.
projection: order_id, customer_id, total_amount, created_at
primary-keys: order_id
description: Filter by order status and prune columns
pipeline:
name: MySQL to SelectDB Pipeline
Full PostgreSQL database synchronization
Synchronize all tables in a PostgreSQL schema to SelectDB. Before you start, set wal_level to logical on the PostgreSQL instance and make sure that max_replication_slots and max_wal_senders have enough headroom. For more information, see Configure a PostgreSQL database.
source:
type: postgres
name: PostgreSQL Source
hostname: <yourHostname>
port: 5432
username: ${secret_values.pg_username}
password: ${secret_values.pg_password}
# PostgreSQL table names use the database.schema.table format.
tables: test_db.public.\.*
slot.name: <yourSlotName>
decoding.plugin.name: pgoutput
scan.startup.mode: initial
# (Optional) Use an existing publication instead of creating one when the job starts.
debezium.publication.autocreate.mode: disabled
debezium.publication.name: <yourPublicationName>
sink:
type: doris
name: SelectDB Sink
fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
jdbc-url: jdbc:mysql://selectdb-cn-****.selectdbfe.rds.aliyuncs.com:9030
username: ${secret_values.selectdb_username}
password: ${secret_values.selectdb_password}
pipeline:
name: PostgreSQL to SelectDB Pipeline
Type mapping
Flink to SelectDB
| Flink CDC type | SelectDB type | Notes |
|---|---|---|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
DECIMAL |
DECIMAL |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
BOOLEAN |
BOOLEAN |
|
DATE |
DATE |
|
TIMESTAMP[(p)] |
DATETIME[(p)] |
|
TIMESTAMP_LTZ[(p)] |
DATETIME[(p)] |
|
CHAR(n) |
CHAR(n*3) |
SelectDB stores strings in UTF-8. English characters occupy 1 byte; Chinese characters occupy 3 bytes. Maximum CHAR length is 255; longer values are auto-converted to VARCHAR. |
VARCHAR(n) |
VARCHAR(n*3) |
Same UTF-8 multiplier applies. Maximum VARCHAR length is 65533; longer values are auto-converted to STRING. |
BINARY(n) |
STRING |
|
VARBINARY(n) |
STRING |
|
STRING |
STRING |