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());
+    }
+
+}

Reply via email to