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