goutamadwant opened a new pull request, #12477:
URL: https://github.com/apache/seatunnel/pull/12477

   ### Purpose of this pull request
   
   Since #11060 was merged, 
`OceanBaseCDCCompatibilityIT.testOceanBaseCdcWrapperRuntimeE2e` has been 
failing on Flink in full CI runs on `dev` and on PRs based on it. Zeta passes 
and Spark is disabled for this test. The Flink client fails while it builds the 
job graph:
   
   ```
   Caused by: org.apache.flink.streaming.runtime.tasks.StreamTaskException: 
Could not serialize object for key serializedUDF.
        at 
org.apache.flink.streaming.api.graph.StreamConfig.lambda$serializeAllConfigs$1(StreamConfig.java:198)
        ...
   Caused by: java.io.NotSerializableException: io.debezium.relational.TableId
        at java.base/java.io.ObjectOutputStream.writeObject0(Unknown Source)
        at java.base/java.io.ObjectOutputStream.writeObject(Unknown Source)
        at java.base/java.util.HashMap.internalWriteEntries(Unknown Source)
        at java.base/java.util.HashMap.writeObject(Unknown Source)
        ...
        at 
org.apache.flink.util.InstantiationUtil.serializeObject(InstantiationUtil.java:553)
   ```
   
   Examples:
   
   - `dev` push build on `deb16a3c3`, `all-connectors-it-5 (8, ubuntu-latest)`: 
4 errors out of 5 invocations (the four Flink versions); Zeta passes. 
https://github.com/apache/seatunnel/actions/runs/35999263398/job/107632711731
   - `dev` nightly schedule on the same commit, `all-connectors-it-5`: 
https://github.com/apache/seatunnel/actions/runs/36025021304
   - An unrelated PR, `all-connectors-it-5 (8, ubuntu-latest)` (Flink 1.18 and 
1.20 fail; Zeta passes): 
https://github.com/goutamadwant/seatunnel/actions/runs/36093925266/job/107942628914
   
   **Cause.** This is a packaging problem, not a problem in the connector code. 
`connector-cdc-base` ships a patched `io.debezium.relational.TableId` that 
implements `Serializable`, along with a few other patched Debezium classes. CDC 
sources rely on this: for example, `AbstractDebeziumDeserializationSchema` has 
a `HashMap<TableId, byte[]>` that Flink serializes with the source. 
`connector-cdc-mysql` excludes `debezium-core` and `debezium-api` so that its 
jar does not carry another copy. Those exclusions are declared in the 
`dependencyManagement` of `connector-cdc-mysql`. `connector-cdc-oceanbase` gets 
MySQL CDC as a transitive dependency, and those exclusions are not applied 
there. As a result, the shaded OceanBase jar also bundles `debezium-core` and 
`debezium-api` 1.9.8, including the original, non-serializable `TableId`. The 
build log already shows this:
   
   ```
   [INFO] Including io.debezium:debezium-api:jar:1.9.8.Final in the shaded jar.
   [INFO] Including io.debezium:debezium-core:jar:1.9.8.Final in the shaded jar.
   ```
   
   This happens in reactor builds, which is how CI and `seatunnel-dist` build 
the connector. A standalone build of the module resolves the installed 
`connector-cdc-mysql` pom, which the shade plugin has already reduced, so the 
leak does not appear there.
   
   Plugin discovery always puts `connector-cdc-base` next to a CDC connector 
jar, so the Flink client sees two `TableId` classes. When the OceanBase jar 
comes first, Flink loads the unpatched class, and serializing the source fails 
as shown above.
   
   **Change.** The `connector-cdc-mysql` dependency in 
`connector-cdc-oceanbase/pom.xml` now excludes `io.debezium:debezium-core` and 
`io.debezium:debezium-api`, as `connector-cdc-mysql` itself does. After this 
change, the OceanBase jar contains the same entries as the MySQL CDC jar, plus 
the two OceanBase classes and their Maven metadata. Its size drops from 28.6 MB 
to 22.3 MB.
   
   `OceanBaseCdcPackagingIT` is added. It runs in the module's failsafe phase 
after packaging, following the existing packaging ITs such as 
`TestPayPalPackagingIT`. The failsafe configuration passes it the path of the 
module's own packaged jar. It checks two things:
   
   1. The packaged jar contains no `io/debezium/**` class that 
`connector-cdc-base` also provides.
   2. A class loader with the OceanBase jar before the `connector-cdc-base` jar 
(the failing order) loads a `TableId` that can be Java-serialized in a 
`HashMap`. This is the same operation that fails in Flink.
   
   ### Does this PR introduce _any_ user-facing change?
   
   Yes, it fixes a bug in the `OceanBase-CDC` source jar. Only `dev` is 
affected, because the connector was added in #11060 and has not been released 
yet. On Flink, jobs with this source could fail at submission with 
`NotSerializableException: io.debezium.relational.TableId`. Options and 
behavior are unchanged.
   
   ### How was this patch tested?
   
   **Packaging IT, red then green.**
   
   | Code | JDK 8 | JDK 11 |
   |---|---|---|
   | Unchanged `dev` (`deb16a3c3`) | 2/2 fail: jar bundles 630 Debezium classes 
that cdc-base also provides; `NotSerializableException: 
io.debezium.relational.TableId` | same 2/2 failures |
   | This patch | 2/2 pass | 2/2 pass |
   
   Command: `./mvnw -B verify -pl 
seatunnel-connectors-v2/connector-cdc/connector-cdc-oceanbase -am -DskipUT=true 
-DskipIT=false -Dit.test=OceanBaseCdcPackagingIT 
-Dit.failIfNoSpecifiedTests=false`
   
   **E2E, red then green.** `OceanBaseCDCCompatibilityIT` with 
`RUN_ALL_CONTAINER=false`, which is the PR engine set: Flink 1.18, Flink 1.20 
and Zeta. Spark is disabled by the test.
   
   | Code | Result |
   |---|---|
   | Unchanged `dev` | Flink 1.20 and Zeta passed. Flink 1.18 hit the test's 60 
s snapshot wait: on this machine the amd64 Flink image runs under emulation, 
and the job took about 110 s to reach RUNNING. This is not the `TableId` error. 
|
   | This patch | 3/3 passed (Flink 1.18, Flink 1.20, Zeta) |
   
   This local run does not show the exact CI failure. 
`AbstractPluginDiscovery.selectPluginJar` orders jars as `File.listFiles()` 
returns them, and that order is not guaranteed. In the CI logs, `Discovery 
plugin jar for: ... OceanBase-CDC` lists the OceanBase jar first, which is the 
failing order. On my machine, `connector-cdc-base` came first. The packaging IT 
above fixes the class loader order so that the failure reproduces the same way 
everywhere. Locally, the E2E run confirms that the source still works on Flink 
1.18/1.20 and Zeta with the slimmer jar.
   
   Command: `TEST_IN_PR=true RUN_ALL_CONTAINER=false ./mvnw -B verify -pl 
:connector-cdc-oceanbase-e2e -am -DskipUT=true -DskipIT=false 
-Dit.test=OceanBaseCDCCompatibilityIT -Dit.failIfNoSpecifiedTests=false` (JDK 8 
host).
   
   **Unit tests.** Unit tests for `connector-cdc-base` (98), 
`connector-cdc-mysql` (33) and `connector-cdc-oceanbase` (6), plus the new 
packaging IT (2), pass on JDK 8 and JDK 11.
   
   **Not verified locally:** Flink 1.13 and 1.15. They fail on `dev` for the 
same reason, and the fix does not depend on the Flink version.
   
   ### 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)
   * [ ] If necessary, please update the documentation to describe the new 
feature. https://github.com/apache/seatunnel/tree/dev/docs
   * [ ] If necessary, please update `incompatible-changes.md` to describe the 
incompatibility caused by this PR.
   * [ ] If you are contributing the connector code, please check that the 
following files are updated:
     1. Update 
[plugin-mapping.properties](https://github.com/apache/seatunnel/blob/dev/plugin-mapping.properties)
 and add new connector information in it
     2. Update the pom file of 
[seatunnel-dist](https://github.com/apache/seatunnel/blob/dev/seatunnel-dist/pom.xml)
     3. Add ci label in 
[label-scope-conf](https://github.com/apache/seatunnel/blob/dev/.github/workflows/labeler/label-scope-conf.yml)
     4. Add e2e testcase in 
[seatunnel-e2e](https://github.com/apache/seatunnel/tree/dev/seatunnel-e2e/seatunnel-connector-v2-e2e/)
     5. Update connector 
[plugin_config](https://github.com/apache/seatunnel/blob/dev/config/plugin_config)
   


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