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]