This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new ac1a0dcc9 [INLONG-7708][Sort] Fix parse source failed when adding a
table during all database migration in Oracle (#7709)
ac1a0dcc9 is described below
commit ac1a0dcc901f0d096b5f243a03aca662374bfe1a
Author: emhui <[email protected]>
AuthorDate: Wed Mar 29 14:22:00 2023 +0800
[INLONG-7708][Sort] Fix parse source failed when adding a table during all
database migration in Oracle (#7709)
---
.../base/relational/JdbcSourceEventDispatcher.java | 39 ++--
.../reader/IncrementalSourceRecordEmitter.java | 2 +-
.../inlong/sort/cdc/base/util/RecordUtils.java | 25 ++
.../base/relational/JdbcSourceEventDispatcher.java | 252 ---------------------
.../oracle/source/reader/OracleRecordEmitter.java | 2 +-
.../reader/fetch/OracleSourceFetchTaskContext.java | 3 +-
.../relational/OracleSourceEventDispatcher.java | 118 ++++++++++
7 files changed, 168 insertions(+), 273 deletions(-)
diff --git
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
index d87929309..0702a2bb4 100644
---
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
+++
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
@@ -17,6 +17,8 @@
package org.apache.inlong.sort.cdc.base.relational;
+import static
org.apache.inlong.sort.cdc.base.util.RecordUtils.isMysqlConnector;
+
import io.debezium.config.CommonConnectorConfig;
import io.debezium.connector.base.ChangeEventQueue;
import io.debezium.document.DocumentWriter;
@@ -60,8 +62,7 @@ import org.slf4j.LoggerFactory;
*/
public class JdbcSourceEventDispatcher extends EventDispatcher<TableId> {
- private static final Logger LOG = LoggerFactory.getLogger(
-
com.ververica.cdc.connectors.base.relational.JdbcSourceEventDispatcher.class);
+ private static final Logger LOG =
LoggerFactory.getLogger(JdbcSourceEventDispatcher.class);
public static final String HISTORY_RECORD_FIELD = "historyRecord";
public static final String SERVER_ID_KEY = "server_id";
@@ -70,14 +71,14 @@ public class JdbcSourceEventDispatcher extends
EventDispatcher<TableId> {
private static final DocumentWriter DOCUMENT_WRITER =
DocumentWriter.defaultWriter();
- private final ChangeEventQueue<DataChangeEvent> queue;
- private final HistorizedDatabaseSchema historizedSchema;
- private final DataCollectionFilters.DataCollectionFilter<TableId> filter;
- private final CommonConnectorConfig connectorConfig;
- private final TopicSelector<TableId> topicSelector;
- private final Schema schemaChangeKeySchema;
- private final Schema schemaChangeValueSchema;
- private final String topic;
+ public final ChangeEventQueue<DataChangeEvent> queue;
+ public final HistorizedDatabaseSchema historizedSchema;
+ public final DataCollectionFilters.DataCollectionFilter<TableId> filter;
+ public final CommonConnectorConfig connectorConfig;
+ public final TopicSelector<TableId> topicSelector;
+ public final Schema schemaChangeKeySchema;
+ public final Schema schemaChangeValueSchema;
+ public final String topic;
public JdbcSourceEventDispatcher(
CommonConnectorConfig connectorConfig,
@@ -181,7 +182,7 @@ public class JdbcSourceEventDispatcher extends
EventDispatcher<TableId> {
}
/** A {@link SchemaChangeEventEmitter.Receiver} implementation for {@link
SchemaChangeEvent}. */
- private final class SchemaChangeEventReceiver implements
SchemaChangeEventEmitter.Receiver {
+ public final class SchemaChangeEventReceiver implements
SchemaChangeEventEmitter.Receiver {
private Struct schemaChangeRecordKey(SchemaChangeEvent event) {
Struct result = new Struct(schemaChangeKeySchema);
@@ -190,14 +191,16 @@ public class JdbcSourceEventDispatcher extends
EventDispatcher<TableId> {
}
private Struct schemaChangeRecordValue(SchemaChangeEvent event) throws
IOException {
- Struct sourceInfo = event.getSource();
Map<String, Object> source = new HashMap<>();
- String fileName = sourceInfo.getString(BINLOG_FILENAME_OFFSET_KEY);
- Long pos = sourceInfo.getInt64(BINLOG_POSITION_OFFSET_KEY);
- Long serverId = sourceInfo.getInt64(SERVER_ID_KEY);
- source.put(SERVER_ID_KEY, serverId);
- source.put(BINLOG_FILENAME_OFFSET_KEY, fileName);
- source.put(BINLOG_POSITION_OFFSET_KEY, pos);
+ if (isMysqlConnector(event.getSource())) {
+ Struct sourceInfo = event.getSource();
+ String fileName =
sourceInfo.getString(BINLOG_FILENAME_OFFSET_KEY);
+ Long pos = sourceInfo.getInt64(BINLOG_POSITION_OFFSET_KEY);
+ Long serverId = sourceInfo.getInt64(SERVER_ID_KEY);
+ source.put(SERVER_ID_KEY, serverId);
+ source.put(BINLOG_FILENAME_OFFSET_KEY, fileName);
+ source.put(BINLOG_POSITION_OFFSET_KEY, pos);
+ }
HistoryRecord historyRecord =
new HistoryRecord(
source,
diff --git
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/reader/IncrementalSourceRecordEmitter.java
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/reader/IncrementalSourceRecordEmitter.java
index faf04aa28..5766eb01e 100644
---
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/reader/IncrementalSourceRecordEmitter.java
+++
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/reader/IncrementalSourceRecordEmitter.java
@@ -23,7 +23,7 @@ import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.getFetch
import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.getHistoryRecord;
import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.getMessageTimestamp;
import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.isDataChangeRecord;
-import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.isSchemaChangeEvent;
+import static
org.apache.inlong.sort.cdc.base.util.RecordUtils.isSchemaChangeEvent;
import io.debezium.document.Array;
import io.debezium.relational.history.HistoryRecord;
diff --git
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/util/RecordUtils.java
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/util/RecordUtils.java
index e859399e5..91ecd992e 100644
---
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/util/RecordUtils.java
+++
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/util/RecordUtils.java
@@ -35,6 +35,7 @@ import org.apache.flink.table.types.logical.SmallIntType;
import org.apache.flink.table.types.logical.TimestampType;
import org.apache.flink.table.types.logical.TinyIntType;
import org.apache.flink.table.types.logical.VarCharType;
+import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.data.Struct;
import org.apache.kafka.connect.source.SourceRecord;
@@ -52,6 +53,10 @@ public class RecordUtils {
.asList("CHAR", "NCHAR", "NVARCHAR2", "NVCHAER", "VARCHAR",
"VARCHAR2", "CLOB", "NCLOB", "XMLType");
private static final List<String> BINARY_TYPE = Arrays.asList("BLOB",
"ROWID");
private static final List<String> BIGINT_TYPE = Arrays.asList("INTERVAL
DAY TO SECOND", "INTERVAL YEAR TO MONTH");
+ public static final String MYSQL_SCHEMA_CHANGE_EVENT_KEY_NAME =
"io.debezium.connector.mysql.SchemaChangeKey";
+ public static final String ORACLE_SCHEMA_CHANGE_EVENT_KEY_NAME =
"io.debezium.connector.oracle.SchemaChangeKey";
+ public static final String CONNECTOR = "connector";
+ public static final String MYSQL_CONNECTOR = "mysql";
private RecordUtils() {
@@ -130,4 +135,24 @@ public class RecordUtils {
return null;
}
+ /**
+ * Whether the source Record is a schema change event.
+ * @param sourceRecord
+ * @return ture if the source Record is a schema change event.
+ */
+ public static boolean isSchemaChangeEvent(SourceRecord sourceRecord) {
+ Schema keySchema = sourceRecord.keySchema();
+ return keySchema != null &&
(MYSQL_SCHEMA_CHANGE_EVENT_KEY_NAME.equalsIgnoreCase(keySchema.name())
+ ||
ORACLE_SCHEMA_CHANGE_EVENT_KEY_NAME.equalsIgnoreCase(keySchema.name()));
+ }
+
+ /**
+ * Whether the source belong Mysql Connector
+ * @param source
+ * @return true if the source belong Mysql Connector
+ */
+ public static boolean isMysqlConnector(Struct source) {
+ String connector = source.getString(CONNECTOR);
+ return MYSQL_CONNECTOR.equalsIgnoreCase(connector);
+ }
}
diff --git
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
deleted file mode 100644
index cb2dd4a60..000000000
---
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/relational/JdbcSourceEventDispatcher.java
+++ /dev/null
@@ -1,252 +0,0 @@
-/*
- * 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.inlong.sort.cdc.base.relational;
-
-import io.debezium.config.CommonConnectorConfig;
-import io.debezium.connector.base.ChangeEventQueue;
-import io.debezium.document.DocumentWriter;
-import io.debezium.pipeline.DataChangeEvent;
-import io.debezium.pipeline.EventDispatcher;
-import io.debezium.pipeline.source.spi.EventMetadataProvider;
-import io.debezium.pipeline.spi.ChangeEventCreator;
-import io.debezium.pipeline.spi.SchemaChangeEventEmitter;
-import io.debezium.relational.TableId;
-import io.debezium.relational.history.HistoryRecord;
-import io.debezium.schema.DataCollectionFilters;
-import io.debezium.schema.DatabaseSchema;
-import io.debezium.schema.HistorizedDatabaseSchema;
-import io.debezium.schema.SchemaChangeEvent;
-import io.debezium.schema.TopicSelector;
-import io.debezium.util.SchemaNameAdjuster;
-import java.io.IOException;
-import java.util.Collection;
-import java.util.HashMap;
-import java.util.Map;
-import org.apache.inlong.sort.cdc.base.source.meta.offset.Offset;
-import org.apache.inlong.sort.cdc.base.source.meta.split.SourceSplitBase;
-import org.apache.inlong.sort.cdc.base.source.meta.wartermark.WatermarkEvent;
-import org.apache.inlong.sort.cdc.base.source.meta.wartermark.WatermarkKind;
-import org.apache.kafka.connect.data.Schema;
-import org.apache.kafka.connect.data.SchemaBuilder;
-import org.apache.kafka.connect.data.Struct;
-import org.apache.kafka.connect.source.SourceRecord;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * A subclass implementation of {@link EventDispatcher}.
- *
- * <pre>
- * 1. This class shares one {@link ChangeEventQueue} between multiple readers.
- * 2. This class override some methods for dispatching {@link HistoryRecord}
directly,
- * this is useful for downstream to deserialize the {@link HistoryRecord}
back.
- * </pre>
- * This class is a rewrite of JdbcSourceEventDispatcher in cdc-base,
- * because there is a conflict between the io-debezium-core:1.5.4-final
depended on in cdc-base
- * and the io-debezium-core:1.6.4-final relied on in this module.
- * Copy from com.ververica:flink-cdc-base:2.3.0.
- */
-public class JdbcSourceEventDispatcher extends EventDispatcher<TableId> {
-
- private static final Logger LOG = LoggerFactory.getLogger(
-
com.ververica.cdc.connectors.base.relational.JdbcSourceEventDispatcher.class);
-
- public static final String HISTORY_RECORD_FIELD = "historyRecord";
- public static final String SERVER_ID_KEY = "server_id";
- public static final String BINLOG_FILENAME_OFFSET_KEY = "file";
- public static final String BINLOG_POSITION_OFFSET_KEY = "pos";
-
- private static final DocumentWriter DOCUMENT_WRITER =
DocumentWriter.defaultWriter();
-
- private final ChangeEventQueue<DataChangeEvent> queue;
- private final HistorizedDatabaseSchema historizedSchema;
- private final DataCollectionFilters.DataCollectionFilter<TableId> filter;
- private final CommonConnectorConfig connectorConfig;
- private final TopicSelector<TableId> topicSelector;
- private final Schema schemaChangeKeySchema;
- private final Schema schemaChangeValueSchema;
- private final String topic;
-
- public JdbcSourceEventDispatcher(
- CommonConnectorConfig connectorConfig,
- TopicSelector<TableId> topicSelector,
- DatabaseSchema<TableId> schema,
- ChangeEventQueue<DataChangeEvent> queue,
- DataCollectionFilters.DataCollectionFilter<TableId> filter,
- ChangeEventCreator changeEventCreator,
- EventMetadataProvider metadataProvider,
- SchemaNameAdjuster schemaNameAdjuster) {
- super(
- connectorConfig,
- topicSelector,
- schema,
- queue,
- filter,
- changeEventCreator,
- metadataProvider,
- schemaNameAdjuster);
- this.historizedSchema =
- schema instanceof HistorizedDatabaseSchema
- ? (HistorizedDatabaseSchema<TableId>) schema
- : null;
- this.filter = filter;
- this.queue = queue;
- this.connectorConfig = connectorConfig;
- this.topicSelector = topicSelector;
- this.topic = topicSelector.getPrimaryTopic();
- this.schemaChangeKeySchema =
- SchemaBuilder.struct()
- .name(
- schemaNameAdjuster.adjust(
- "io.debezium.connector."
- +
connectorConfig.getConnectorName()
- + ".SchemaChangeKey"))
- .field(HistoryRecord.Fields.DATABASE_NAME,
Schema.STRING_SCHEMA)
- .build();
- this.schemaChangeValueSchema =
- SchemaBuilder.struct()
- .name(
- schemaNameAdjuster.adjust(
- "io.debezium.connector."
- +
connectorConfig.getConnectorName()
- + ".SchemaChangeValue"))
- .field(
- HistoryRecord.Fields.SOURCE,
-
connectorConfig.getSourceInfoStructMaker().schema())
- .field(HISTORY_RECORD_FIELD,
Schema.OPTIONAL_STRING_SCHEMA)
- .build();
- }
-
- public ChangeEventQueue<DataChangeEvent> getQueue() {
- return queue;
- }
-
- @Override
- public void dispatchSchemaChangeEvent(
- TableId dataCollectionId, SchemaChangeEventEmitter
schemaChangeEventEmitter)
- throws InterruptedException {
- if (dataCollectionId != null && !filter.isIncluded(dataCollectionId)) {
- if (historizedSchema == null ||
historizedSchema.storeOnlyCapturedTables()) {
- LOG.trace("Filtering schema change event for {}",
dataCollectionId);
- return;
- }
- }
- schemaChangeEventEmitter.emitSchemaChangeEvent(new
SchemaChangeEventReceiver());
- }
-
- @Override
- public void dispatchSchemaChangeEvent(
- Collection<TableId> dataCollectionIds,
- SchemaChangeEventEmitter schemaChangeEventEmitter)
- throws InterruptedException {
- boolean anyNonfilteredEvent = false;
- if (dataCollectionIds == null || dataCollectionIds.isEmpty()) {
- anyNonfilteredEvent = true;
- } else {
- for (TableId dataCollectionId : dataCollectionIds) {
- if (filter.isIncluded(dataCollectionId)) {
- anyNonfilteredEvent = true;
- break;
- }
- }
- }
- if (!anyNonfilteredEvent) {
- if (historizedSchema == null ||
historizedSchema.storeOnlyCapturedTables()) {
- LOG.trace("Filtering schema change event for {}",
dataCollectionIds);
- return;
- }
- }
-
- schemaChangeEventEmitter.emitSchemaChangeEvent(new
SchemaChangeEventReceiver());
- }
-
- /** A {@link SchemaChangeEventEmitter.Receiver} implementation for {@link
SchemaChangeEvent}. */
- private final class SchemaChangeEventReceiver implements
SchemaChangeEventEmitter.Receiver {
-
- private Struct schemaChangeRecordKey(SchemaChangeEvent event) {
- Struct result = new Struct(schemaChangeKeySchema);
- result.put(HistoryRecord.Fields.DATABASE_NAME,
event.getDatabase());
- return result;
- }
-
- private Struct schemaChangeRecordValue(SchemaChangeEvent event) throws
IOException {
- Struct sourceInfo = event.getSource();
- Map<String, Object> source = new HashMap<>();
- String fileName = sourceInfo.getString(BINLOG_FILENAME_OFFSET_KEY);
- Long pos = sourceInfo.getInt64(BINLOG_POSITION_OFFSET_KEY);
- Long serverId = sourceInfo.getInt64(SERVER_ID_KEY);
- source.put(SERVER_ID_KEY, serverId);
- source.put(BINLOG_FILENAME_OFFSET_KEY, fileName);
- source.put(BINLOG_POSITION_OFFSET_KEY, pos);
- HistoryRecord historyRecord =
- new HistoryRecord(
- source,
- event.getOffset(),
- event.getDatabase(),
- null,
- event.getDdl(),
- event.getTableChanges());
- String historyStr =
DOCUMENT_WRITER.write(historyRecord.document());
-
- Struct value = new Struct(schemaChangeValueSchema);
- value.put(HistoryRecord.Fields.SOURCE, event.getSource());
- value.put(HISTORY_RECORD_FIELD, historyStr);
- return value;
- }
-
- @Override
- public void schemaChangeEvent(SchemaChangeEvent event) throws
InterruptedException {
- historizedSchema.applySchemaChange(event);
- if (connectorConfig.isSchemaChangesHistoryEnabled()) {
- try {
- final String topicName = topicSelector.getPrimaryTopic();
- final Integer partition = 0;
- final Struct key = schemaChangeRecordKey(event);
- final Struct value = schemaChangeRecordValue(event);
- final SourceRecord record =
- new SourceRecord(
- event.getPartition(),
- event.getOffset(),
- topicName,
- partition,
- schemaChangeKeySchema,
- key,
- schemaChangeValueSchema,
- value);
- queue.enqueue(new DataChangeEvent(record));
- } catch (IOException e) {
- throw new IllegalStateException(
- String.format("dispatch schema change event %s
error ", event), e);
- }
- }
- }
- }
-
- public void dispatchWatermarkEvent(
- Map<String, ?> sourcePartition,
- SourceSplitBase sourceSplit,
- Offset watermark,
- WatermarkKind watermarkKind)
- throws InterruptedException {
-
- SourceRecord sourceRecord =
- WatermarkEvent.create(
- sourcePartition, topic, sourceSplit.splitId(),
watermarkKind, watermark);
- queue.enqueue(new DataChangeEvent(sourceRecord));
- }
-}
diff --git
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/OracleRecordEmitter.java
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/OracleRecordEmitter.java
index df7e0d2ee..610229903 100644
---
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/OracleRecordEmitter.java
+++
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/OracleRecordEmitter.java
@@ -21,7 +21,7 @@ import static
com.ververica.cdc.connectors.base.source.meta.wartermark.Watermark
import static
com.ververica.cdc.connectors.base.source.meta.wartermark.WatermarkEvent.isWatermarkEvent;
import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.getHistoryRecord;
import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.isDataChangeRecord;
-import static
com.ververica.cdc.connectors.base.utils.SourceRecordUtils.isSchemaChangeEvent;
+import static
org.apache.inlong.sort.cdc.base.util.RecordUtils.isSchemaChangeEvent;
import com.ververica.cdc.debezium.history.FlinkJsonTableChangeSerializer;
import io.debezium.connector.AbstractSourceInfo;
diff --git
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/fetch/OracleSourceFetchTaskContext.java
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/fetch/OracleSourceFetchTaskContext.java
index 93656a223..d4919dd36 100644
---
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/fetch/OracleSourceFetchTaskContext.java
+++
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/reader/fetch/OracleSourceFetchTaskContext.java
@@ -53,6 +53,7 @@ import
org.apache.inlong.sort.cdc.base.source.meta.split.SourceSplitBase;
import
org.apache.inlong.sort.cdc.base.source.reader.external.JdbcSourceFetchTaskContext;
import org.apache.inlong.sort.cdc.oracle.source.config.OracleSourceConfig;
import org.apache.inlong.sort.cdc.oracle.source.meta.offset.RedoLogOffset;
+import
org.apache.inlong.sort.cdc.oracle.source.relational.OracleSourceEventDispatcher;
import org.apache.inlong.sort.cdc.oracle.source.utils.OracleUtils;
import org.apache.kafka.connect.data.Struct;
import org.apache.kafka.connect.source.SourceRecord;
@@ -123,7 +124,7 @@ public class OracleSourceFetchTaskContext extends
JdbcSourceFetchTaskContext {
// .buffering()
.build();
this.dispatcher =
- new JdbcSourceEventDispatcher(
+ new OracleSourceEventDispatcher(
connectorConfig,
topicSelector,
databaseSchema,
diff --git
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/relational/OracleSourceEventDispatcher.java
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/relational/OracleSourceEventDispatcher.java
new file mode 100644
index 000000000..43da1915a
--- /dev/null
+++
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/relational/OracleSourceEventDispatcher.java
@@ -0,0 +1,118 @@
+/*
+ * 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.inlong.sort.cdc.oracle.source.relational;
+
+import io.debezium.config.CommonConnectorConfig;
+import io.debezium.connector.base.ChangeEventQueue;
+import io.debezium.document.DocumentWriter;
+import io.debezium.pipeline.DataChangeEvent;
+import io.debezium.pipeline.EventDispatcher;
+import io.debezium.pipeline.source.spi.EventMetadataProvider;
+import io.debezium.pipeline.spi.ChangeEventCreator;
+import io.debezium.pipeline.spi.SchemaChangeEventEmitter;
+import io.debezium.relational.TableId;
+import io.debezium.relational.history.HistoryRecord;
+import io.debezium.schema.DataCollectionFilters;
+import io.debezium.schema.DatabaseSchema;
+import io.debezium.schema.TopicSelector;
+import io.debezium.util.SchemaNameAdjuster;
+import java.util.Collection;
+import org.apache.inlong.sort.cdc.base.relational.JdbcSourceEventDispatcher;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * A subclass implementation of {@link EventDispatcher}.
+ *
+ * <pre>
+ * 1. This class shares one {@link ChangeEventQueue} between multiple readers.
+ * 2. This class override some methods for dispatching {@link HistoryRecord}
directly,
+ * this is useful for downstream to deserialize the {@link HistoryRecord}
back.
+ * </pre>
+ * This class extends JdbcSourceEventDispatcher in cdc-base,
+ * because there is a conflict between the io-debezium-core:1.5.4-final
depended on in cdc-base
+ * and the io-debezium-core:1.6.4-final relied on in this module.
+ */
+public class OracleSourceEventDispatcher extends JdbcSourceEventDispatcher {
+
+ private static final Logger LOG =
LoggerFactory.getLogger(OracleSourceEventDispatcher.class);
+
+ public static final String HISTORY_RECORD_FIELD = "historyRecord";
+
+ private static final DocumentWriter DOCUMENT_WRITER =
DocumentWriter.defaultWriter();
+
+ public OracleSourceEventDispatcher(
+ CommonConnectorConfig connectorConfig,
+ TopicSelector<TableId> topicSelector,
+ DatabaseSchema<TableId> schema,
+ ChangeEventQueue<DataChangeEvent> queue,
+ DataCollectionFilters.DataCollectionFilter<TableId> filter,
+ ChangeEventCreator changeEventCreator,
+ EventMetadataProvider metadataProvider,
+ SchemaNameAdjuster schemaNameAdjuster) {
+ super(
+ connectorConfig,
+ topicSelector,
+ schema,
+ queue,
+ filter,
+ changeEventCreator,
+ metadataProvider,
+ schemaNameAdjuster);
+ }
+
+ @Override
+ public void dispatchSchemaChangeEvent(
+ TableId dataCollectionId, SchemaChangeEventEmitter
schemaChangeEventEmitter)
+ throws InterruptedException {
+ if (dataCollectionId != null && !filter.isIncluded(dataCollectionId)) {
+ if (historizedSchema == null ||
historizedSchema.storeOnlyCapturedTables()) {
+ LOG.trace("Filtering schema change event for {}",
dataCollectionId);
+ return;
+ }
+ }
+ schemaChangeEventEmitter.emitSchemaChangeEvent(new
SchemaChangeEventReceiver());
+ }
+
+ @Override
+ public void dispatchSchemaChangeEvent(
+ Collection<TableId> dataCollectionIds,
+ SchemaChangeEventEmitter schemaChangeEventEmitter)
+ throws InterruptedException {
+ boolean anyNonfilteredEvent = false;
+ if (dataCollectionIds == null || dataCollectionIds.isEmpty()) {
+ anyNonfilteredEvent = true;
+ } else {
+ for (TableId dataCollectionId : dataCollectionIds) {
+ if (filter.isIncluded(dataCollectionId)) {
+ anyNonfilteredEvent = true;
+ break;
+ }
+ }
+ }
+ if (!anyNonfilteredEvent) {
+ if (historizedSchema == null ||
historizedSchema.storeOnlyCapturedTables()) {
+ LOG.trace("Filtering schema change event for {}",
dataCollectionIds);
+ return;
+ }
+ }
+
+ schemaChangeEventEmitter.emitSchemaChangeEvent(new
SchemaChangeEventReceiver());
+ }
+
+}