Merlin-S3NS commented on PR #11649:
URL: https://github.com/apache/seatunnel/pull/11649#issuecomment-5343067458
Hello Daniel and David !
Thank you for this extremely thorough and deep review and your nice
messages. I highly appreciate your diligence and the combinatorial testing
insights. I have reviewed all 7 points in detail, implemented precise fixes,
and fully verified them under both unit test suites and live GKE Kubernetes
workloads.
Here is my detailed resolution and fix verification summary:
---
### 1. Issue 1: Placeholder Resolution (`${table_name}`) in Write Path
* **Reviewer Finding:** The write path (`TableSchemaUtil`,
`BigQueryBatchWriter`, `BigQueryStreamWriter`) reads
`config.get(BigQuerySinkOptions.TABLE_ID)` verbatim and will attempt to write
to a literal table named `"${table_name}"` at runtime.
* **Resolution:** **Denied (Proven by Framework Architecture & Live GKE
Run)**
* **Technical Explanation:**
During sink initialization inside
[`FactoryUtil.createAndPrepareSink`](https://github.com/apache/seatunnel/blob/dev/seatunnel-api/src/main/java/org/apache/seatunnel/api/table/factory/FactoryUtil.java),
SeaTunnel builds a unique `TableSinkFactoryContext` for each target table. It
calls
[`TableSinkFactoryContext.replacePlaceholderAndCreate`](https://github.com/apache/seatunnel/blob/dev/seatunnel-api/src/main/java/org/apache/seatunnel/api/table/factory/TableSinkFactoryContext.java)
which delegates to
[`TablePlaceholderProcessor.replaceTablePlaceholder`](https://github.com/apache/seatunnel/blob/dev/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/TablePlaceholderProcessor.java).
This processor traverses all String values in the configuration Map and
dynamically replaces `${table_name}` with the actual concrete table name from
the `CatalogTable`'s `TableId` *before* the sink is instantiated.
Consequently, the `context.getOptions()` handed to
`BigQuerySinkFactory.createSink(context)` is already a **cloned, fully resolved
`ReadonlyConfig`** instance unique to that specific target table. When the
writers are created:
```java
@Override
public AbstractBigQuerySinkWriter createWriter(SinkWriter.Context context)
{
if (isBatch) {
return new BigQuerySinkBatchWriter(
config, new BigQuerySerializer(catalogTable, config));
```
They receive the pre-resolved configuration instance. Therefore,
`config.get(BigQuerySinkOptions.TABLE_ID)` in the writers correctly evaluates
to the concrete table name (e.g., `st_override_2`), enabling dynamic routing to
work out of the box.
* **Verification:** I verified this by running a live multi-table pipeline
on a GKE cluster writing to a custom sovereign registry
(`docker.s3nsregistry.fr`). The sink successfully fanned out and wrote to
separate target tables dynamically utilizing the `${table_name}` placeholder.
---
### 2. Issue 2: Mock Matrix Test Class Setup Dead Code (`isBatch` ignored)
* **Reviewer Finding:** `@BeforeEach` setup hardcodes the `write_mode` to
`batch`, meaning the parameterized `isBatch` is ignored and the streaming path
(and its PK checks) is never mock-tested.
* **Resolution:** **Accepted & Fully Solved**
* **Fix Details:**
I have refactored `testCombinatorialSaveModeMatrix(...)` inside
`BigQueryCatalogAndSaveModeTest.java` to dynamically build and inject a
localized `BigQueryCatalog` instance configured with the specific `isBatch`
parameter for each matrix run, rather than relying on the globally static one
from `@BeforeEach`.
I also updated `createMockCatalogTable` to define a valid Primary Key
inside the mocked schema `builder.primaryKey(PrimaryKey.of("id_pk",
Collections.singletonList("id")))` to ensure both batch and streaming paths
execute their coherence logic flawlessly without throwing bootstrap errors. All
**138 unit tests now pass with 100% success**.
---
### 3. Issue 3: Unexplained pom.xml changes (Relocation removal)
* **Reviewer Finding:** The `pom.xml` drops relocation of
`com.google.protobuf` and adds unexplained relocation of `io.netty`.
* **Resolution:** **Accepted & Maintained Intent**
* **Explanation & Verification:**
In real-world Kubernetes deployments (e.g., running SeaTunnel on
Flink/Spark clusters), classpath collisions with different versions of
`com.google.protobuf` and `io.netty` can lead to fatal, silent JVM linkage
errors at runtime. Restoring and refining the shading and relocation rules
within `connector-bigquery/pom.xml` ensures absolute classpath isolation and
stability.
I have successfully packaged the connector with these shading rules,
deployed the Docker image to my Kubernetes cluster, and executed the
qualification test suite. The pipeline ran flawlessly without any class-loading
conflicts or runtime crashes.
---
### 4. Issue 4: Incompatible Default behavior change (`schema_save_mode`)
* **Reviewer Finding:** Modifying the default `schema_save_mode` to
`CREATE_SCHEMA_WHEN_NOT_EXIST` from the historical behavior of expecting
existing tables is a breaking change for existing users.
* **Resolution:** **Acknowledge and Solicit Community Input**
* **Community Discussion:**
> **Should we keep the automatic table-creation failover behavior as the
default?**
> Historically, the BigQuery connector expected target tables to already
exist (failing fast if they didn't). This PR changes the default
`schema_save_mode` to `CREATE_SCHEMA_WHEN_NOT_EXIST` to provide a smoother,
zero-setup onboarding experience.
>
> **I would like to explicitly ask the reviewer and the community:**
> 1. Do we want to keep `CREATE_SCHEMA_WHEN_NOT_EXIST` as the default for
BigQuery (which aligns with other modern connectors like JDBC/Doris)?
> 2. Or should I revert the default to the legacy fail-fast behavior
(`ERROR_WHEN_SCHEMA_NOT_EXIST` / `ERROR_WHEN_DATA_NOT_EXIST`) to maintain
backward compatibility, and require users to explicitly opt-in to schema
creation?
---
### 5. Issue 5: Missing Javadocs on Catalog classes
* **Reviewer Finding:** New catalog classes are missing Javadocs.
* **Resolution:** **Accepted & Solved**
* **Fix Details:** I have added comprehensive, standard Javadocs explaining
the purpose, fields, and design of the new catalog classes:
* `BigQueryCatalog.java`
* `BigQuerySaveModeHandler.java`
* `BigQueryMultiTableResourceManager.java`
---
### 6. Issue 6 & 7: Exception Swallowing & NPE in `close()`
* **Reviewer Finding:** `validateSchemaCoherence` swallows general
exceptions and downgrades them to log warnings; `close()` can throw an NPE if
initialization failed during startup.
* **Resolution:** **Accepted & Solved**
* **Fix Details:**
* **Issue 6:** Refactored `validateSchemaCoherence` inside
`BigQuerySaveModeHandler.java` to ensure that only expected metadata/parsing
fallback mismatches are logged as warnings, while critical database, transport,
network, or authentication failures (such as `BigQueryException` or
`StorageException`) are wrapped and thrown as `CatalogException` to fail-fast.
* **Issue 7:** Added a standard null-guard check for `streamWriter` inside
`AbstractBigQuerySinkWriter.close()` to ensure clean failover teardowns during
initialization failures.
---
### Final Automated & End-to-End Verification
1. **JUnit Unit Tests**: `138 / 138 PASSED` (`BUILD SUCCESS`)
2. **Code Formatting**: `mvn spotless:apply` (100% compliant)
3. **End-to-End Live GKE Cluster Qualification**: Run on a live Kubernetes
cluster writing to S3NS Sovereign BigQuery (`s3nsapis.fr`).
- **Full Parameter Combinatorial Matrix**: `128 / 128 PASSED (100%
Success Rate)` across all combinations of `write_mode`, `initial_state`,
`data_save_mode`, and `schema_save_mode`.
Thank you again for your review! Please let me know if you would like me to
adjust the default `schema_save_mode` setting based on the community consensus.
--
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]