Merlin-S3NS opened a new pull request, #11649:
URL: https://github.com/apache/seatunnel/pull/11649
<!--
Thank you for contributing to SeaTunnel! Please make sure that your code
changes
are covered with tests. And in case of new features or big changes
remember to adjust the documentation.
Feel free to ping committers for the review!
## Contribution Checklist
- Make sure that the pull request corresponds to a [GITHUB
issue](https://github.com/apache/seatunnel/issues).
- Name the pull request in the form "[Feature] [component] Title of the
pull request", where *Feature* can be replaced by `Hotfix`, `Bug`, etc.
- Minor fixes should be named following this pattern: `[hotfix] [docs] Fix
typo in README.md doc`.
-->
### Purpose of this pull request
This pull request implements comprehensive **SaveMode (automatic DDL/schema
synchronization)** and **Multi-Table write support** for the Google Cloud
BigQuery sink connector.
Specifically, it implements:
1. **`SupportSaveMode`**: Integrated `BigQuerySaveModeHandler` extending the
default handler to support schema-save modes (`CREATE_SCHEMA_WHEN_NOT_EXIST`,
`RECREATE_SCHEMA`, `ERROR_WHEN_SCHEMA_NOT_EXIST`, `IGNORE`) and data-save modes
(`APPEND_DATA`, `DROP_DATA`, `CUSTOM_PROCESSING`, `ERROR_WHEN_DATA_EXISTS`).
2. **Schema Coherence Validation**: Implemented robust type compatibility
and column checking to validate the coherence between the upstream schema and
target BigQuery tables during initialization.
3. **`SupportMultiTableSink`**: Implemented multi-table coordination via
`BigQueryMultiTableResourceManager` to gracefully share high-performance
`BigQueryWriteClient` instances across multiple writing tasks, enabling
seamless concurrent multi-table pipelines.
4. **Dynamic Table Routing**: Supported dynamic `${table_name}`-like string
placeholders in `table_id` and dataset_id` settings to dynamically partition
and route multiple upstream tables into distinct BigQuery tables.
5. **Detailed Documentation**: Updated English and Chinese connector
reference guides with option descriptions, tables, dynamic routing examples,
and save mode properties.
---
### Does this PR introduce _any_ user-facing change?
**Yes.** This PR introduces new configuration properties and behavior
upgrades for the BigQuery connector.
#### 1. New Configuration Parameters
Three optional SaveMode options have been exposed for the BigQuery sink:
- `schema_save_mode` (default: `CREATE_SCHEMA_WHEN_NOT_EXIST`): Controls the
target schema preparation before sync.
- `data_save_mode` (default: `APPEND_DATA`): Controls how existing target
table data is treated.
- `custom_sql`: Specifies a custom SQL statement to execute if
`data_save_mode` is `CUSTOM_PROCESSING`.
#### 2. Upgraded Table Setup Behavior
- **Previous Behavior**: Target BigQuery tables were required to exist prior
to job submission. Upstream multi-table routing required splitting pipelines
with multiple independent sink blocks.
- **New Behavior**: Target BigQuery tables and datasets are automatically
created at runtime using metadata from the upstream schemas if
`schema_save_mode` is set to `CREATE_SCHEMA_WHEN_NOT_EXIST` or
`RECREATE_SCHEMA`. Multiple tables can be dynamically routed to corresponding
BigQuery tables by specifying `table_id = "${table_name}"`.
---
### How was this patch tested?
This patch has been rigorously verified through both comprehensive unit
tests and high-concurrency real-world qualification tests in a GKE environment.
#### 1. Local Unit & Metadata Tests
A new test class `BigQueryCatalogAndSaveModeTest` was introduced, verifying
all execution matrices for schema save modes, data save modes, database and
table auto-creation, column compatibility validation, type widening checks, and
mock client integrations.
Run command:
```bash
export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64
./mvnw test -pl seatunnel-connectors-v2/connector-bigquery
```
##### Unit Test Execution Summary Table
| Test Class | Tests Run | Failures | Errors | Skipped | Time |
| :--- | :---: | :---: | :---: | :---: | :---: |
| **`BigQueryCatalogAndSaveModeTest`** | **104** | **0** | **0** | **0** |
**1.511 s** |
| `BigQuerySinkBatchWriterTest` | 4 | 0 | 0 | 0 | 0.017 s |
| `BigQuerySinkFactoryTest` | 2 | 0 | 0 | 0 | 0.398 s |
| `TableSchemaUtilTest` | 2 | 0 | 0 | 0 | 0.040 s |
| `BigQuerySinkRestoreStateTest` | 1 | 0 | 0 | 0 | 0.001 s |
| `BigQuerySerializerTest` | 24 | 0 | 0 | 0 | 0.043 s |
| **Total** | **137** | **0** | **0** | **0** | **2.010 s** |
**Status:** `BUILD SUCCESS` (137/137 tests successfully passed on Java 11).
---
#### 2. Real-World GKE Qualification Matrix Tests
To guarantee the highest level of production readiness and compatibility, an
exhaustive parallel GKE test was orchestrated using the `run_csv_matrix.py`
test suite.
##### GKE Test Environment & Orchestration:
- **Orchestration Platform**: SeaTunnel Cluster running in a Google
Kubernetes Engine (GKE) namespace `seatunnel` (Scaled to 3 masters and 10
workers for parallel execution).
- **Active Master Pods**: 3
- **Target BigQuery Cloud**: Sovereign Cloud GKE environment (S3NS Sandbox)
targeting `universe_domain = "s3nsapis.fr"`.
##### Test Generation & Verification Strategy:
The `run_csv_matrix.py` test engine dynamically compiled and deployed **128
combinatorial HOCON test cases** in parallel across the GKE masters. Each case
evaluated a unique combination of:
- **Write Modes**: `BATCH`, `STREAMING`
- **Initial Target States**: `NOT_EXIST` (No dataset/table), `EMPTY_TABLE`,
`EXIST_1_ROW`, `EXIST_WRONG_SCHEMA`
- **Data Save Modes**: `APPEND_DATA`, `DROP_DATA`, `ERROR_WHEN_DATA_EXISTS`,
`CUSTOM_PROCESSING`
- **Schema Save Modes**: `RECREATE_SCHEMA`, `CREATE_SCHEMA_WHEN_NOT_EXIST`,
`ERROR_WHEN_SCHEMA_NOT_EXIST`, `IGNORE`
##### Matrix Results:
- **Total Tests Executed**: `128`
- **Passed Cases**: `128`
- **Failed Cases**: `0`
All 128 combinatorial states successfully passed verification. This confirms
full production readiness, high concurrency stability, and seamless SaveMode
operation in Kubernetes sovereign deployments.
---
### 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)
(N/A - no new Jar packages added)
* [x] If necessary, please update the documentation to describe the new
feature. https://github.com/apache/seatunnel/tree/dev/docs (English and Chinese
BigQuery connector sink documents fully refreshed inline)
* [ ] If necessary, please update `incompatible-changes.md` to describe the
incompatibility caused by this PR. (N/A - fully backward-compatible feature
addition)
* [x] If you are contributing the connector code, please check that the
following files are updated:
- *Note: BigQuery is an existing connector in SeaTunnel, so
`plugin-mapping.properties`, `pom` files, sidebars, and ci labels are already
fully configured. This PR updates documentation files `BigQuery.md` (English &
Chinese) inline.*
* [ ] If you are contributing the engine code, please check that the
following files are updated: (N/A)
[128_matrix_actual_results.csv](https://github.com/user-attachments/files/30701538/128_matrix_actual_results.csv)
[128_matrix_results.html](https://github.com/user-attachments/files/30701540/128_matrix_results.html)
[128_matrix_expected_results.csv](https://github.com/user-attachments/files/30701541/128_matrix_expected_results.csv)
--
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]