This is an automated email from the ASF dual-hosted git repository.

pacinogong 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 21262a4d4 [INLONG-7639][Sort] Metric label list constants are issued 
by each connector (#7640)
21262a4d4 is described below

commit 21262a4d4377a6e8a3b382285ebfd780fa4e98c4
Author: emhui <[email protected]>
AuthorDate: Tue Mar 21 10:09:11 2023 +0800

    [INLONG-7639][Sort] Metric label list constants are issued by each 
connector (#7640)
---
 .../inlong/sort/cdc/base/config/MetricConfig.java  |   9 +
 .../sort/cdc/base/source/IncrementalSource.java    |   2 +-
 .../sort/cdc/base/source/IncrementalSource.java    | 230 ---------------------
 .../oracle/source/config/OracleSourceConfig.java   |   7 +
 4 files changed, 17 insertions(+), 231 deletions(-)

diff --git 
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/config/MetricConfig.java
 
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/config/MetricConfig.java
index 33a48af4a..5c1facc0d 100644
--- 
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/config/MetricConfig.java
+++ 
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/config/MetricConfig.java
@@ -18,6 +18,7 @@
 package org.apache.inlong.sort.cdc.base.config;
 
 import java.io.Serializable;
+import java.util.List;
 
 /** The mertic configuration which offers basic metric configuration. **/
 public interface MetricConfig extends Serializable {
@@ -36,4 +37,12 @@ public interface MetricConfig extends Serializable {
      */
     String getInlongAudit();
 
+    /**
+     * getMetricLabelList
+     *
+     * @return metric label list of each connector.
+     * eg: oracle metric label list is [DATABASE, SCHEMA, TABLE]
+     */
+    List<String> getMetricLabelList();
+
 }
diff --git 
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
 
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
index 8278cc7cf..39969d208 100644
--- 
a/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
+++ 
b/inlong-sort/sort-connectors/cdc-base/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
@@ -126,7 +126,7 @@ public class IncrementalSource<T, C extends SourceConfig>
                 .withRegisterMetric(RegisteredMetric.ALL)
                 .build();
 
-        sourceReaderMetrics.registerMetrics(metricOption);
+        sourceReaderMetrics.registerMetrics(metricOption, 
metricConfig.getMetricLabelList());
         Supplier<IncrementalSourceSplitReader<C>> splitReaderSupplier =
                 () -> new IncrementalSourceSplitReader<>(
                         readerContext.getIndexOfSubtask(), dataSourceDialect, 
sourceConfig);
diff --git 
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
 
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
deleted file mode 100644
index 30650ced5..000000000
--- 
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/base/source/IncrementalSource.java
+++ /dev/null
@@ -1,230 +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.source;
-
-import com.ververica.cdc.connectors.base.options.StartupMode;
-import io.debezium.relational.TableId;
-import java.lang.reflect.Method;
-import java.util.Arrays;
-import java.util.List;
-import java.util.function.Supplier;
-import org.apache.flink.annotation.Experimental;
-import org.apache.flink.api.common.typeinfo.TypeInformation;
-import org.apache.flink.api.connector.source.Boundedness;
-import org.apache.flink.api.connector.source.Source;
-import org.apache.flink.api.connector.source.SourceReaderContext;
-import org.apache.flink.api.connector.source.SplitEnumerator;
-import org.apache.flink.api.connector.source.SplitEnumeratorContext;
-import org.apache.flink.api.java.typeutils.ResultTypeQueryable;
-import org.apache.flink.connector.base.source.reader.RecordEmitter;
-import org.apache.flink.connector.base.source.reader.RecordsWithSplitIds;
-import 
org.apache.flink.connector.base.source.reader.synchronization.FutureCompletingBlockingQueue;
-import org.apache.flink.core.io.SimpleVersionedSerializer;
-import org.apache.flink.metrics.MetricGroup;
-import org.apache.flink.util.FlinkRuntimeException;
-import org.apache.inlong.sort.base.Constants;
-import org.apache.inlong.sort.base.metric.MetricOption;
-import org.apache.inlong.sort.base.metric.MetricOption.RegisteredMetric;
-import org.apache.inlong.sort.cdc.base.config.MetricConfig;
-import org.apache.inlong.sort.cdc.base.config.SourceConfig;
-import org.apache.inlong.sort.cdc.base.debezium.DebeziumDeserializationSchema;
-import org.apache.inlong.sort.cdc.base.dialect.DataSourceDialect;
-import org.apache.inlong.sort.cdc.base.source.assigner.HybridSplitAssigner;
-import org.apache.inlong.sort.cdc.base.source.assigner.SplitAssigner;
-import org.apache.inlong.sort.cdc.base.source.assigner.StreamSplitAssigner;
-import 
org.apache.inlong.sort.cdc.base.source.assigner.state.HybridPendingSplitsState;
-import 
org.apache.inlong.sort.cdc.base.source.assigner.state.PendingSplitsState;
-import 
org.apache.inlong.sort.cdc.base.source.assigner.state.PendingSplitsStateSerializer;
-import 
org.apache.inlong.sort.cdc.base.source.assigner.state.StreamPendingSplitsState;
-import 
org.apache.inlong.sort.cdc.base.source.enumerator.IncrementalSourceEnumerator;
-import org.apache.inlong.sort.cdc.base.source.meta.offset.OffsetFactory;
-import org.apache.inlong.sort.cdc.base.source.meta.split.SourceRecords;
-import org.apache.inlong.sort.cdc.base.source.meta.split.SourceSplitBase;
-import org.apache.inlong.sort.cdc.base.source.meta.split.SourceSplitSerializer;
-import org.apache.inlong.sort.cdc.base.source.meta.split.SourceSplitState;
-import org.apache.inlong.sort.cdc.base.source.metrics.SourceReaderMetrics;
-import org.apache.inlong.sort.cdc.base.source.reader.IncrementalSourceReader;
-import 
org.apache.inlong.sort.cdc.base.source.reader.IncrementalSourceRecordEmitter;
-import 
org.apache.inlong.sort.cdc.base.source.reader.IncrementalSourceSplitReader;
-
-/**
- * The basic source of Incremental Snapshot framework for datasource, it is 
based on FLIP-27 and
- * Watermark Signal Algorithm which supports parallel reading snapshot of 
table and then continue to
- * capture data change by streaming reading.
- * Copy from com.ververica:flink-cdc-base:2.3.0.
- */
-@Experimental
-public class IncrementalSource<T, C extends SourceConfig>
-        implements
-            Source<T, SourceSplitBase, PendingSplitsState>,
-            ResultTypeQueryable<T> {
-
-    private static final long serialVersionUID = 1L;
-
-    protected final SourceConfig.Factory<C> configFactory;
-    protected final DataSourceDialect<C> dataSourceDialect;
-    protected final OffsetFactory offsetFactory;
-    protected final DebeziumDeserializationSchema<T> deserializationSchema;
-    protected final SourceSplitSerializer sourceSplitSerializer;
-
-    public IncrementalSource(
-            SourceConfig.Factory<C> configFactory,
-            DebeziumDeserializationSchema<T> deserializationSchema,
-            OffsetFactory offsetFactory,
-            DataSourceDialect<C> dataSourceDialect) {
-        this.configFactory = configFactory;
-        this.deserializationSchema = deserializationSchema;
-        this.offsetFactory = offsetFactory;
-        this.dataSourceDialect = dataSourceDialect;
-        this.sourceSplitSerializer =
-                new SourceSplitSerializer() {
-
-                    @Override
-                    public OffsetFactory getOffsetFactory() {
-                        return offsetFactory;
-                    }
-                };
-    }
-
-    @Override
-    public Boundedness getBoundedness() {
-        return Boundedness.CONTINUOUS_UNBOUNDED;
-    }
-
-    @Override
-    public IncrementalSourceReader<T, C> createReader(SourceReaderContext 
readerContext)
-            throws Exception {
-        // create source config for the given subtask (e.g. unique server id)
-        C sourceConfig = 
configFactory.create(readerContext.getIndexOfSubtask());
-        MetricConfig metricConfig = (MetricConfig) sourceConfig;
-        FutureCompletingBlockingQueue<RecordsWithSplitIds<SourceRecords>> 
elementsQueue =
-                new FutureCompletingBlockingQueue<>();
-
-        // Forward compatible with flink 1.13
-        final Method metricGroupMethod = 
readerContext.getClass().getMethod("metricGroup");
-        metricGroupMethod.setAccessible(true);
-        final MetricGroup metricGroup = (MetricGroup) 
metricGroupMethod.invoke(readerContext);
-        final SourceReaderMetrics sourceReaderMetrics = new 
SourceReaderMetrics(metricGroup);
-
-        // create source config for the given subtask (e.g. unique server id)
-        MetricOption metricOption = MetricOption.builder()
-                .withInlongLabels(metricConfig.getInlongMetric())
-                .withAuditAddress(metricConfig.getInlongAudit())
-                .withRegisterMetric(RegisteredMetric.ALL)
-                .build();
-
-        sourceReaderMetrics.registerMetrics(metricOption,
-                Arrays.asList(Constants.DATABASE_NAME, Constants.SCHEMA_NAME, 
Constants.TABLE_NAME));
-        Supplier<IncrementalSourceSplitReader<C>> splitReaderSupplier =
-                () -> new IncrementalSourceSplitReader<>(
-                        readerContext.getIndexOfSubtask(), dataSourceDialect, 
sourceConfig);
-        return new IncrementalSourceReader<>(
-                elementsQueue,
-                splitReaderSupplier,
-                createRecordEmitter(sourceConfig, sourceReaderMetrics),
-                readerContext.getConfiguration(),
-                readerContext,
-                sourceConfig,
-                sourceSplitSerializer,
-                dataSourceDialect,
-                sourceReaderMetrics);
-    }
-
-    @Override
-    public SplitEnumerator<SourceSplitBase, PendingSplitsState> 
createEnumerator(
-            SplitEnumeratorContext<SourceSplitBase> enumContext) {
-        C sourceConfig = configFactory.create(0);
-        final SplitAssigner splitAssigner;
-        if (sourceConfig.getStartupOptions().startupMode == 
StartupMode.INITIAL) {
-            try {
-                final List<TableId> remainingTables =
-                        
dataSourceDialect.discoverDataCollections(sourceConfig);
-                boolean isTableIdCaseSensitive =
-                        
dataSourceDialect.isDataCollectionIdCaseSensitive(sourceConfig);
-                splitAssigner =
-                        new HybridSplitAssigner<>(
-                                sourceConfig,
-                                enumContext.currentParallelism(),
-                                remainingTables,
-                                isTableIdCaseSensitive,
-                                dataSourceDialect,
-                                offsetFactory);
-            } catch (Exception e) {
-                throw new FlinkRuntimeException(
-                        "Failed to discover captured tables for enumerator", 
e);
-            }
-        } else {
-            splitAssigner = new StreamSplitAssigner(sourceConfig, 
dataSourceDialect, offsetFactory);
-        }
-
-        return new IncrementalSourceEnumerator(enumContext, sourceConfig, 
splitAssigner);
-    }
-
-    @Override
-    public SplitEnumerator<SourceSplitBase, PendingSplitsState> 
restoreEnumerator(
-            SplitEnumeratorContext<SourceSplitBase> enumContext, 
PendingSplitsState checkpoint) {
-        C sourceConfig = configFactory.create(0);
-
-        final SplitAssigner splitAssigner;
-        if (checkpoint instanceof HybridPendingSplitsState) {
-            splitAssigner =
-                    new HybridSplitAssigner<>(
-                            sourceConfig,
-                            enumContext.currentParallelism(),
-                            (HybridPendingSplitsState) checkpoint,
-                            dataSourceDialect,
-                            offsetFactory);
-        } else if (checkpoint instanceof StreamPendingSplitsState) {
-            splitAssigner =
-                    new StreamSplitAssigner(
-                            sourceConfig,
-                            (StreamPendingSplitsState) checkpoint,
-                            dataSourceDialect,
-                            offsetFactory);
-        } else {
-            throw new UnsupportedOperationException(
-                    "Unsupported restored PendingSplitsState: " + checkpoint);
-        }
-        return new IncrementalSourceEnumerator(enumContext, sourceConfig, 
splitAssigner);
-    }
-
-    @Override
-    public SimpleVersionedSerializer<SourceSplitBase> getSplitSerializer() {
-        return sourceSplitSerializer;
-    }
-
-    @Override
-    public SimpleVersionedSerializer<PendingSplitsState> 
getEnumeratorCheckpointSerializer() {
-        SourceSplitSerializer sourceSplitSerializer = (SourceSplitSerializer) 
getSplitSerializer();
-        return new PendingSplitsStateSerializer(sourceSplitSerializer);
-    }
-
-    @Override
-    public TypeInformation<T> getProducedType() {
-        return deserializationSchema.getProducedType();
-    }
-
-    protected RecordEmitter<SourceRecords, T, SourceSplitState> 
createRecordEmitter(
-            SourceConfig sourceConfig, SourceReaderMetrics 
sourceReaderMetrics) {
-        return new IncrementalSourceRecordEmitter<>(
-                deserializationSchema,
-                sourceReaderMetrics,
-                sourceConfig.isIncludeSchemaChanges(),
-                offsetFactory);
-    }
-}
diff --git 
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/config/OracleSourceConfig.java
 
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/config/OracleSourceConfig.java
index a033f1e44..7d668719a 100644
--- 
a/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/config/OracleSourceConfig.java
+++ 
b/inlong-sort/sort-connectors/oracle-cdc/src/main/java/org/apache/inlong/sort/cdc/oracle/source/config/OracleSourceConfig.java
@@ -22,9 +22,11 @@ import io.debezium.config.Configuration;
 import io.debezium.connector.oracle.OracleConnectorConfig;
 import io.debezium.relational.RelationalTableFilters;
 import java.time.Duration;
+import java.util.Arrays;
 import java.util.List;
 import java.util.Properties;
 import javax.annotation.Nullable;
+import org.apache.inlong.sort.base.Constants;
 import org.apache.inlong.sort.cdc.base.config.JdbcSourceConfig;
 
 /**
@@ -108,4 +110,9 @@ public class OracleSourceConfig extends JdbcSourceConfig {
     public String getUrl() {
         return url;
     }
+
+    @Override
+    public List<String> getMetricLabelList() {
+        return Arrays.asList(Constants.DATABASE_NAME, Constants.SCHEMA_NAME, 
Constants.TABLE_NAME);
+    }
 }

Reply via email to