zhangshenghang commented on code in PR #11060:
URL: https://github.com/apache/seatunnel/pull/11060#discussion_r3844971398


##########
seatunnel-connectors-v2/connector-cdc/connector-cdc-oceanbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oceanbase/source/OceanBaseIncrementalSourceFactory.java:
##########
@@ -0,0 +1,132 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.cdc.oceanbase.source;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.connector.TableSource;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceTableConfig;
+import org.apache.seatunnel.connectors.cdc.base.option.JdbcSourceOptions;
+import org.apache.seatunnel.connectors.cdc.base.option.SourceOptions;
+import org.apache.seatunnel.connectors.cdc.base.utils.CatalogTableUtils;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.config.MySqlSourceConfigFactory;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.MySqlIncrementalSourceFactory;
+
+import com.google.auto.service.AutoService;
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.Serializable;
+import java.util.List;
+import java.util.Optional;
+
+/**
+ * Factory for the OceanBase CDC source.
+ *
+ * <p>The first implementation deliberately wraps the MySQL CDC connector 
because that is the stable
+ * and testable path already adopted by Flink CDC for OceanBase Binlog Service.
+ */
+@AutoService(Factory.class)
+@Slf4j
+public class OceanBaseIncrementalSourceFactory extends 
MySqlIncrementalSourceFactory {
+
+    /**
+     * Return the identifier used in SeaTunnel source config.
+     *
+     * @return OceanBase CDC identifier
+     */
+    @Override
+    public String factoryIdentifier() {
+        return OceanBaseIncrementalSource.IDENTIFIER;
+    }
+
+    /**
+     * Return the concrete source class used for OceanBase CDC jobs.
+     *
+     * @return OceanBase CDC source class
+     */
+    @Override
+    public Class<? extends SeaTunnelSource> getSourceClass() {
+        return OceanBaseIncrementalSource.class;
+    }
+
+    /**
+     * Restore the source by reusing MySQL-compatible table discovery and 
checkpoint merge logic.
+     *
+     * @param context source factory context
+     * @param restoreTables restored table structures from checkpoint state
+     * @param <T> emitted record type
+     * @param <SplitT> split type
+     * @param <StateT> state type
+     * @return restorable OceanBase CDC table source
+     */
+    @Override
+    public <T, SplitT extends SourceSplit, StateT extends Serializable>
+            TableSource<T, SplitT, StateT> restoreSource(
+                    TableSourceFactoryContext context, List<CatalogTable> 
restoreTables) {
+        return () -> {
+            try {
+                Class.forName("com.mysql.cj.jdbc.Driver");
+            } catch (Exception e) {
+                log.warn("Failed to load JDBC driver com.mysql.cj.jdbc.Driver 
", e);
+            }
+            ReadonlyConfig config = context.getOptions();
+            List<CatalogTable> catalogTables =
+                    CatalogTableUtil.getCatalogTables(config, 
context.getClassLoader());

Review Comment:
   **[P1] The documented option contract is broken: this call requires 
`compatible_mode` at runtime, but neither the docs nor `optionRule()` declare 
it**
   
   The chain: `CatalogTableUtil.getCatalogTables(config, classLoader)` derives 
the catalog factory id from `plugin_name.replace("-CDC", "")` 
(CatalogTableUtil.java:94), so `OceanBase-CDC` resolves catalog factory 
`"OceanBase"` → `OceanBaseCatalogFactory` from connector-jdbc (bundled 
transitively via connector-cdc-mysql). `FactoryUtil.createOptionalCatalog` 
(FactoryUtil.java:297) then validates `OceanBaseCatalogFactory.optionRule()`, 
which is 
`JdbcCommonOptions.baseCatalogRule().required(JdbcCommonOptions.COMPATIBLE_MODE)`.
   
   I verified this against the current head with a config that matches the doc 
task example exactly (no `compatible_mode`):
   
   ```
   OptionValidationException: Option validation failed (1 error):
     [1] option: 'compatible_mode'
         type: required
         constraint: required option is not configured
        at ConfigValidator.validate(ConfigValidator.java:183)
        at FactoryUtil.createOptionalCatalog(FactoryUtil.java:295)
   ```
   
   Consequences:
   
   1. Both task examples (`docs/en/connectors/source/OceanBase-CDC.md` and the 
zh counterpart) omit `compatible_mode`, so a user copying them hits this 
failure immediately, with nothing in the docs explaining the key.
   2. The docs state "`OceanBase-CDC` intentionally reuses the same option 
contract as `MySQL-CDC`" — that is not true at runtime: one extra required key 
exists.
   3. The contract is also inconsistent in the other direction: because 
`optionRule()` (inherited from `MySqlIncrementalSourceFactory`) does not 
declare `compatible_mode`, the static validation path 
(`SeaTunnelConfValidateCommand` → `ConfigValidator.validateUnknownKeys`) 
rejects configs that *do* set it — including this PR's own E2E conf — with 
`unknown option keys: [compatible_mode]`. The E2E only passes because the 
runtime path does not check unknown keys.
   4. The value is unconstrained: `compatible_mode = "oracle"` routes snapshot 
schema discovery to `OceanBaseOracleCatalog` while the incremental runtime 
stays MySQL binlog — a combination this PR explicitly declares out of scope.
   
   Suggested fix (any one of):
   
   - Pin MySQL mode in the wrapper: default `compatible_mode` to `"mysql"` in 
`restoreSource` before the catalog call (rebuild the `ReadonlyConfig` from 
`getSourceMap()` with the key added when absent) — this keeps the "same 
contract as MySQL-CDC" claim true; or
   - Declare it: override `optionRule()` to add 
`JdbcCommonOptions.COMPATIBLE_MODE` and document the key in both doc files and 
both task examples; or
   - Call `CatalogTableUtil.getCatalogTables("MySQL", config, classLoader)` 
explicitly so the MySQL catalog is used.
   
   Please also add a unit test that exercises `restoreSource` (or at least the 
catalog discovery) without `compatible_mode` — the current 
`OceanBaseIncrementalSourceFactoryTest` only asserts 
identifier/inheritance/class, which is why this wasn't caught in earlier review 
rounds.



##########
seatunnel-connectors-v2/connector-cdc/connector-cdc-oceanbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oceanbase/source/OceanBaseIncrementalSourceFactory.java:
##########
@@ -0,0 +1,132 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.cdc.oceanbase.source;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.connector.TableSource;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceTableConfig;
+import org.apache.seatunnel.connectors.cdc.base.option.JdbcSourceOptions;
+import org.apache.seatunnel.connectors.cdc.base.option.SourceOptions;
+import org.apache.seatunnel.connectors.cdc.base.utils.CatalogTableUtils;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.config.MySqlSourceConfigFactory;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.MySqlIncrementalSourceFactory;
+
+import com.google.auto.service.AutoService;
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.Serializable;
+import java.util.List;
+import java.util.Optional;
+
+/**
+ * Factory for the OceanBase CDC source.
+ *
+ * <p>The first implementation deliberately wraps the MySQL CDC connector 
because that is the stable
+ * and testable path already adopted by Flink CDC for OceanBase Binlog Service.
+ */
+@AutoService(Factory.class)
+@Slf4j
+public class OceanBaseIncrementalSourceFactory extends 
MySqlIncrementalSourceFactory {
+
+    /**
+     * Return the identifier used in SeaTunnel source config.
+     *
+     * @return OceanBase CDC identifier
+     */
+    @Override
+    public String factoryIdentifier() {
+        return OceanBaseIncrementalSource.IDENTIFIER;
+    }
+
+    /**
+     * Return the concrete source class used for OceanBase CDC jobs.
+     *
+     * @return OceanBase CDC source class
+     */
+    @Override
+    public Class<? extends SeaTunnelSource> getSourceClass() {
+        return OceanBaseIncrementalSource.class;
+    }
+
+    /**
+     * Restore the source by reusing MySQL-compatible table discovery and 
checkpoint merge logic.
+     *
+     * @param context source factory context
+     * @param restoreTables restored table structures from checkpoint state
+     * @param <T> emitted record type
+     * @param <SplitT> split type
+     * @param <StateT> state type
+     * @return restorable OceanBase CDC table source
+     */
+    @Override
+    public <T, SplitT extends SourceSplit, StateT extends Serializable>

Review Comment:
   **[P2] `restoreSource` duplicates the whole body of 
`MySqlIncrementalSourceFactory.restoreSource` — drift risk**
   
   The full method body (driver load, the schema-change fallback chain 
including the `TODO remove this after all users used the new schema change 
option` comment, and the `table-names-config` merge) is copied from 
MySqlIncrementalSourceFactory.java:113-164. Any future change on the MySQL side 
will silently not apply here, and the two connectors' restore semantics can 
drift with no test catching it. That path is actively evolving (the 
schema-change merge and the TODO fallback both landed recently).
   
   Suggestion: extract the shared table-building part into a `protected` method 
on `MySqlIncrementalSourceFactory`, e.g. `protected List<CatalogTable> 
buildCatalogTables(TableSourceFactoryContext context, List<CatalogTable> 
restoreTables)`, and have both factories share it; this class then only 
overrides source construction. That would reduce this factory to the two 
identity methods plus one line, and keeps the restore logic single-sourced.



##########
seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-oceanbase-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oceanbase/OceanBaseCDCCompatibilityIT.java:
##########
@@ -0,0 +1,288 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.cdc.oceanbase;
+
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.EngineType;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
+import org.apache.seatunnel.e2e.common.junit.TestContainerExtension;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.lifecycle.Startables;
+import org.testcontainers.utility.DockerLoggerFactory;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Stream;
+
+import static org.awaitility.Awaitility.await;
+
+/**
+ * Verifies that the OceanBase CDC connector can be discovered and run through 
the MySQL-compatible
+ * CDC path used by the current OceanBase wrapper.
+ *
+ * <p>This is intentionally a compatibility smoke E2E. It does not claim to 
validate a real
+ * OceanBase Binlog Service deployment because the project has no reusable 
OceanBase CDC
+ * Testcontainers environment yet.
+ */
+@Slf4j
+@DisabledOnContainer(
+        value = {},
+        type = {EngineType.SPARK},
+        disabledReason = "Currently SPARK do not support cdc")
+public class OceanBaseCDCCompatibilityIT extends TestSuiteBase implements 
TestResource {
+
+    /**
+     * Network alias used by SeaTunnel engine containers to reach the 
MySQL-compatible CDC source.
+     */
+    private static final String MYSQL_HOST = "mysql_cdc_e2e";
+
+    /** Administrative user used by the test to prepare source and sink 
tables. */
+    private static final String MYSQL_USER_NAME = "mysqluser";
+
+    /** Password shared by the setup SQL users in the MySQL-compatible E2E 
environment. */
+    private static final String MYSQL_USER_PASSWORD = "mysqlpw";
+
+    /** Database name shared with the SeaTunnel job config under test. */
+    private static final String MYSQL_DATABASE = "mysql_cdc";
+
+    /** Small source table used to keep the OceanBase wrapper smoke test 
focused and fast. */
+    private static final String SOURCE_TABLE = 
"oceanbase_cdc_e2e_source_table";
+
+    /** Dedicated sink table generated by the JDBC sink for this compatibility 
test. */
+    private static final String SINK_TABLE = "oceanbase_cdc_e2e_sink_table";
+
+    /** MySQL driver jar required by the OceanBase wrapper because it 
delegates to MySQL CDC. */
+    private static final String MYSQL_DRIVER_URL =
+            
"https://repo1.maven.org/maven2/com/mysql/mysql-connector-j/8.0.32/mysql-connector-j-8.0.32.jar";;
+
+    /** MySQL-compatible source container used as the reproducible CDC runtime 
for the wrapper. */
+    private static final MySqlCompatibleContainer MYSQL_CONTAINER = 
createMySqlContainer();
+
+    /**
+     * Build the MySQL-compatible source with GTID/binlog settings required by 
CDC readers.
+     *
+     * @return configured MySQL-compatible container
+     */
+    private static MySqlCompatibleContainer createMySqlContainer() {
+        return new MySqlCompatibleContainer("8.0.43")
+                .withConfigurationOverride("docker/server-gtids/my.cnf")
+                .withSetupSQL("docker/setup.sql")
+                .withNetwork(NETWORK)
+                .withNetworkAliases(MYSQL_HOST)
+                .withDatabaseName(MYSQL_DATABASE)
+                .withUsername(MYSQL_USER_NAME)
+                .withPassword(MYSQL_USER_PASSWORD)
+                .withLogConsumer(
+                        new 
Slf4jLogConsumer(DockerLoggerFactory.getLogger("mysql-docker-image")));
+    }
+
+    /**
+     * Install the JDBC driver into the OceanBase CDC plugin directory inside 
each engine container.
+     */
+    @TestContainerExtension
+    protected final ContainerExtendedFactory extendedFactory =
+            container -> {
+                Container.ExecResult extraCommands =
+                        container.execInContainer(
+                                "bash",
+                                "-c",
+                                "mkdir -p 
/tmp/seatunnel/plugins/OceanBase-CDC/lib && cd 
/tmp/seatunnel/plugins/OceanBase-CDC/lib && wget "

Review Comment:
   **[P3] Two non-blocking notes on the E2E**
   
   1. This is a wrapper smoke test against stock `mysql:8.0.43` — nothing 
OceanBase-specific is exercised (no OceanBase types, no Binlog Service 
protocol). The class Javadoc discloses this honestly, which is good. Since the 
user docs claim "Requires OceanBase Binlog Service", please consider opening a 
follow-up issue for a real OceanBase CE container test so that claim eventually 
gets validated.
   2. The driver is `wget`-ed from repo1 with a hardcoded `8.0.32` version. The 
newer MySQL CDC E2E pattern (`MysqlCDCDriverResolver`) copies the driver from 
the test classpath instead — that avoids the runtime network dependency and 
keeps the driver version in sync with the pom. Worth aligning when convenient.



-- 
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]

Reply via email to