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]
