DanielLeens commented on code in PR #11060: URL: https://github.com/apache/seatunnel/pull/11060#discussion_r3852086724
########## 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: I re-verified this claim independently against the current head and it holds up completely — thank you for catching it. Confirmed chain, matching the actual source: - `OceanBaseIncrementalSourceFactory` does not override `optionRule()`, so it inherits `MySqlIncrementalSourceFactory.optionRule()` verbatim (MySqlIncrementalSourceFactory.java:56-105) — no `compatible_mode` key anywhere in it. - `restoreSource` calls `CatalogTableUtil.getCatalogTables(config, classLoader)`, which derives `factoryId = plugin_name.replace("-CDC", "")` (CatalogTableUtil.java:94-96) → `"OceanBase"` for this connector. - That routes to `OceanBaseCatalogFactory`, whose `optionRule()` is `JdbcCommonOptions.baseCatalogRule().required(JdbcCommonOptions.COMPATIBLE_MODE)` (OceanBaseCatalogFactory.java:64-68), and `createCatalog()` reads `options.get(JdbcCommonOptions.COMPATIBLE_MODE)` unconditionally (line 44) to pick `OceanBaseOracleCatalog` vs `OceanBaseMySqlCatalog`. So both directions you describe are real: 1. Both doc task examples (en/zh `OceanBase-CDC.md`) omit `compatible_mode` — I checked, it's not present in either — so copy-pasting the doc example throws `OptionValidationException` at catalog-open time. 2. The E2E's own conf (`oceanbasecdc_to_mysql_compatible.conf:24-30`) *does* set `compatible_mode = "mysql"`, but since `optionRule()` doesn't declare it, static `--dry-run` / unknown-key validation would reject that same conf. The E2E only passes because it runs the job for real rather than validating first. 3. And you're right that the value is functionally unconstrained today: `compatible_mode = "oracle"` would send catalog/schema discovery to `OceanBaseOracleCatalog` while the incremental runtime stays plain MySQL binlog (inherited unmodified from `MySqlIncrementalSource`) — a combination the docs explicitly say is out of scope but nothing in code actually prevents. This went through several of my own review rounds without me catching it, since I was checking the wrapper's runtime/checkpoint/restore reuse and treating "same option contract as MySQL-CDC" as a given rather than re-deriving the catalog-factory-id path myself. Apologies for missing it earlier — this is a genuine carryover I should have caught. Of your three suggested fixes, I'd lean toward pinning `compatible_mode = "mysql"` inside `restoreSource` before the catalog call (rebuild `ReadonlyConfig` with the key defaulted when absent) plus declaring it in `optionRule()` — that keeps the "same contract as MySQL-CDC" doc claim true for the common case while also making the static validator accept the E2E's own conf. Declaring it without pinning still leaves the oracle-mode catalog/runtime mismatch open. Agreed this should block merge until resolved, and the added unit test for `restoreSource` without `compatible_mode` is the right regression guard so it can't silently regress again. ########## 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: Confirmed — I diffed the two methods and this is a full copy of `MySqlIncrementalSourceFactory.restoreSource` (MySqlIncrementalSourceFactory.java:113-164) into `OceanBaseIncrementalSourceFactory.restoreSource` (OceanBaseIncrementalSourceFactory.java:82-132), down to the driver-load try/catch and the `TODO remove this after all users used the new schema change option` comment. Only the final constructed type differs (`OceanBaseIncrementalSource` vs `MySqlIncrementalSource`). Agreed this is drift risk rather than a merge blocker — the schema-change fallback/merge logic has changed more than once already in this PR's own review history, and there's no test today that would catch the two factories silently diverging. Your suggested extraction (a `protected buildCatalogTables(...)` on the base factory, shared by both) is the right shape: it keeps `OceanBaseIncrementalSourceFactory` down to the two identity overrides + one call. I'd file this as a non-blocking recommended follow-up rather than something that has to land in this PR, since it's a refactor of the shared MySQL factory, not new OceanBase-specific logic — but I'd support doing it in this PR too if the author has bandwidth, since it directly reduces the surface that P1's fix has to touch. -- 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]
