luozihen opened a new pull request, #12015:
URL: https://github.com/apache/seatunnel/pull/12015
Purpose of this pull request
Close #11878.
Add a new JDBC sink option `multi-table_config.primary_keys` to support
per-table
primary key mapping in multi-table generated-SQL scenarios.
When the upstream table name matches one of the configured patterns, the
mapped key
columns are used; otherwise it falls back to the existing `primary_keys` /
catalog
metadata logic. Inside this new option, `${primary_key}` / `${unique_key}`
can be
mixed with static columns. The legacy top-level `primary_keys` contract is
kept
unchanged.
This PR also fixes the config-parsing issue that blocked the feature:
SeaTunnel
reconstructs config via `ConfigFactory.parseMap(...)`, which interprets map
keys as
Typesafe Config path expressions. Regex keys such as `^t_nova_.*$` contain
`$` and
`.`, so they fail with `ConfigException$BadPath`. The reconstruction now
uses a JSON
round-trip (`ConfigFactory.parseString(..., JSON)`) in
`ConfigShadeUtils.processConfig` and `ReadonlyConfig.toConfig`.
### Does this PR introduce _any_ user-facing change?
Yes. Adds a new optional JDBC sink option `multi-table_config`.
Semantics:
- Each key is a Java regular expression matched against the upstream table
name with
full-match semantics (`tableName.matches(pattern)`).
- Each value is a list of key columns. `${primary_key}` and `${unique_key}`
are
supported only inside this option and can be mixed with static columns.
- Precedence: matched mapping wins; otherwise falls back to the existing
logic
(top-level `primary_keys` -> catalog primary key -> first unique key ->
plain INSERT).
- If a table matches multiple patterns, the first pattern in declaration
order wins.
- If a matched table uses `${primary_key}` / `${unique_key}` but the
upstream table has
no primary/unique key, the job fails with a clear error.
No existing option is renamed, removed, or changed in default behavior.
### Example configurations
#### Example 1
```hocon
env {
job.mode = "STREAMING"
parallelism = 1
checkpoint.interval = 10000
}
source {
MySQL-CDC {
plugin_output = "sharding00_source"
url = "jdbc:mysql://<host>:3306/st_source"
username = "<username>"
password = "<password>"
database-names = ["st_source"]
table-pattern =
"st_source\\.t_(nova_(fo_serial|order)_(000|048)|tyuen_txn_(cp|qr)_(000|048))"
startup.mode = "initial"
server-id = "5400-5408"
}
}
transform {
Sql {
plugin_input = "sharding00_source"
plugin_output = "sharding00_transform"
query = "SELECT *, 'idc' AS DATA_SOURCE FROM dual"
}
}
sink {
Jdbc {
plugin_input = "sharding00_transform"
url =
"jdbc:mysql://<host>:3306/st_target?rewriteBatchedStatements=true&autoReconnect=true"
driver = "com.mysql.cj.jdbc.Driver"
username = "<username>"
password = "<password>"
generate_sink_sql = true
database = "st_target"
multi-table_config {
primary_keys {
"^t_nova_.*$" = ["${primary_key}", "DATA_SOURCE"]
"^t_tyuen_txn_.*$" = ["id_txn_ctrl", "DATA_SOURCE"]
}
}
schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
data_save_mode = "APPEND_DATA"
}
}
```
#### Example 2
```hocon
env {
job.mode = "STREAMING"
parallelism = 1
checkpoint.interval = 10000
}
source {
MySQL-CDC {
plugin_output = "gote_source"
url = "jdbc:mysql://<host>:3306/st_source"
username = "<username>"
password = "<password>"
database-names = ["st_source"]
table-pattern =
"st_source\\.t_(nova_merchant_info|tyuen_txn_ext|nova_merge_settle_serial)"
startup.mode = "initial"
server-id = "5500-5508"
}
}
transform {
Sql {
plugin_input = "gote_source"
plugin_output = "gote_transform"
query = "SELECT *, 'idc' AS DATA_SOURCE FROM dual"
table_transform = [
{
table_path = "st_source.t_nova_merchant_info"
query = "SELECT * FROM dual"
}
]
}
}
sink {
Jdbc {
plugin_input = "gote_transform"
url =
"jdbc:mysql://<host>:3306/st_target?rewriteBatchedStatements=true&autoReconnect=true"
driver = "com.mysql.cj.jdbc.Driver"
username = "<username>"
password = "<password>"
generate_sink_sql = true
database = "st_target"
primary_keys = ["merchant_id"]
multi-table_config {
primary_keys {
"t_tyuen_txn_ext.*" = ["id_txn_ctrl", "DATA_SOURCE"]
"t_nova_merge_settle_serial" = ["${primary_key}", "DATA_SOURCE"]
}
}
schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"
data_save_mode = "APPEND_DATA"
}
}
```
### 建表 SQL
源表和目标表均为手动创建。`multi-table_config.primary_keys` 决定写入目标表时使用的主键列
(体现为最终生成的 upsert / delete 语句的 key 列)。
#### Demo 1 源表
```sql
CREATE TABLE t_nova_fo_serial_000 (
serial_id BIGINT NOT NULL,
order_no VARCHAR(64),
amount DECIMAL(18,2),
status VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (serial_id)
);
CREATE TABLE t_nova_order_000 (
id BIGINT NOT NULL,
order_no VARCHAR(64),
amount DECIMAL(18,2),
created_at DATETIME,
PRIMARY KEY (id)
);
CREATE TABLE t_tyuen_txn_cp_000 (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
amount DECIMAL(18,2),
channel VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (id),
UNIQUE KEY uk_id_txn_ctrl (id_txn_ctrl)
);
CREATE TABLE t_tyuen_txn_qr_000 (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
qr_code VARCHAR(128),
amount DECIMAL(18,2),
created_at DATETIME,
PRIMARY KEY (id),
UNIQUE KEY uk_id_txn_ctrl (id_txn_ctrl)
);
```
#### Demo 1 目标表
```sql
CREATE TABLE t_nova_fo_serial_000 (
serial_id BIGINT NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
order_no VARCHAR(64),
amount DECIMAL(18,2),
status VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (serial_id, DATA_SOURCE)
);
CREATE TABLE t_nova_order_000 (
id BIGINT NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
order_no VARCHAR(64),
amount DECIMAL(18,2),
created_at DATETIME,
PRIMARY KEY (id, DATA_SOURCE)
);
CREATE TABLE t_tyuen_txn_cp_000 (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
amount DECIMAL(18,2),
channel VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (id_txn_ctrl, DATA_SOURCE)
);
CREATE TABLE t_tyuen_txn_qr_000 (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
qr_code VARCHAR(128),
amount DECIMAL(18,2),
created_at DATETIME,
PRIMARY KEY (id_txn_ctrl, DATA_SOURCE)
);
```
#### Demo 2 源表
```sql
CREATE TABLE t_nova_merchant_info (
id BIGINT NOT NULL,
merchant_id VARCHAR(64) NOT NULL,
merchant_name VARCHAR(128),
mcc VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (id),
UNIQUE KEY uk_merchant_id (merchant_id)
);
CREATE TABLE t_tyuen_txn_ext (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
txn_amount DECIMAL(18,2),
ext_info VARCHAR(255),
created_at DATETIME,
PRIMARY KEY (id),
UNIQUE KEY uk_id_txn_ctrl (id_txn_ctrl)
);
CREATE TABLE t_nova_merge_settle_serial (
id BIGINT NOT NULL,
settle_no VARCHAR(64),
amount DECIMAL(18,2),
settle_date DATE,
created_at DATETIME,
PRIMARY KEY (id)
);
```
#### Demo 2 目标表
```sql
CREATE TABLE t_nova_merchant_info (
id BIGINT NOT NULL,
merchant_id VARCHAR(64) NOT NULL,
merchant_name VARCHAR(128),
mcc VARCHAR(16),
created_at DATETIME,
PRIMARY KEY (merchant_id)
);
CREATE TABLE t_tyuen_txn_ext (
id BIGINT NOT NULL,
id_txn_ctrl VARCHAR(64) NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
txn_amount DECIMAL(18,2),
ext_info VARCHAR(255),
created_at DATETIME,
PRIMARY KEY (id_txn_ctrl, DATA_SOURCE)
);
CREATE TABLE t_nova_merge_settle_serial (
id BIGINT NOT NULL,
DATA_SOURCE VARCHAR(32) NOT NULL,
settle_no VARCHAR(64),
amount DECIMAL(18,2),
settle_date DATE,
created_at DATETIME,
PRIMARY KEY (id, DATA_SOURCE)
);
```
### 生成的 SQL 结果
下面是用上面两个 demo 跑通后,MySQL JDBC sink 实际生成的部分 upsert / delete 语句。
可以看出每个目标表已经按 `multi-table_config.primary_keys` 使用了不同的主键列。
```text
-- Example 1
INSERT INTO `st_target`.`t_nova_fo_serial_000`
(`serial_id`, `order_no`, `amount`, `status`, `created_at`, `DATA_SOURCE`)
VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE ...;
DELETE FROM `st_target`.`t_nova_fo_serial_000`
WHERE `serial_id` = ? AND `DATA_SOURCE` = ?;
INSERT INTO `st_target`.`t_tyuen_txn_cp_000`
(`id`, `id_txn_ctrl`, `amount`, `channel`, `created_at`, `DATA_SOURCE`)
VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE ...;
DELETE FROM `st_target`.`t_tyuen_txn_cp_000`
WHERE `id_txn_ctrl` = ? AND `DATA_SOURCE` = ?;
-- Example 2
INSERT INTO `st_target`.`t_nova_merchant_info`
(`id`, `merchant_id`, `merchant_name`, `mcc`, `created_at`)
VALUES (?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE ...;
DELETE FROM `st_target`.`t_nova_merchant_info`
WHERE `merchant_id` = ?;
INSERT INTO `st_target`.`t_tyuen_txn_ext`
(`id`, `id_txn_ctrl`, `txn_amount`, `ext_info`, `created_at`,
`DATA_SOURCE`)
VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE ...;
DELETE FROM `st_target`.`t_tyuen_txn_ext`
WHERE `id_txn_ctrl` = ? AND `DATA_SOURCE` = ?;
INSERT INTO `st_target`.`t_nova_merge_settle_serial`
(`id`, `settle_no`, `amount`, `settle_date`, `created_at`, `DATA_SOURCE`)
VALUES (?, ?, ?, ?, ?, ?) ON DUPLICATE KEY UPDATE ...;
DELETE FROM `st_target`.`t_nova_merge_settle_serial`
WHERE `id` = ? AND `DATA_SOURCE` = ?;
```
### How was this patch tested?
Added unit tests:
- `JdbcSinkFactoryTest`: mapping hit, fallback to legacy `primary_keys`,
missing upstream
primary key error, first-match-wins, string value handling, and the full
factory-context path.
- `ConfigShadeTest`: special-character config keys are preserved after
`decryptConfig`.
- `ReadableConfigTest`: special-character config keys are preserved after
`ReadonlyConfig#toConfig`.
Manually verified with MySQL-CDC -> JDBC using the two demo configurations
above
(`generate_sink_sql = true`), and confirmed the manually-created target
tables together
with the generated upsert / delete statements match the expected per-table
primary key
mapping.
Also run `./mvnw spotless:apply`.
### Check list
* [ ] If any new Jar binary package adding in your PR, please add License
Notice according [New License
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/developer/new-license.md)
* [x] If necessary, please update the documentation to describe the new
feature. https://github.com/apache/seatunnel/tree/dev/docs
* [ ] If necessary, please update `incompatible-changes.md` to describe the
incompatibility caused by this PR.
* [ ] If you are contributing the connector code, please check that the
following files are updated:
1. Update
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
and add new connector information in it
2. Update the pom file of
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
3. Add ci label in
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
4. Add e2e testcase in
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
5. Update connector
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]