listTables(@Nullable String namespace, @Nullable String schemaName) {
+ return Db2SchemaUtils.listTables(sourceConfig, schemaName);
+ }
+
+ /**
+ * Get the {@link Schema} of the given table.
+ *
+ * @param tableId The {@link TableId} of the given table.
+ * @return The {@link Schema} of the table.
+ */
+ @Override
+ public Schema getTableSchema(TableId tableId) {
+ return Db2SchemaUtils.getTableSchema(sourceConfig, tableId);
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineSource.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineSource.java
new file mode 100644
index 00000000000..4814fcc8a34
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineSource.java
@@ -0,0 +1,62 @@
+/*
+ * 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.flink.cdc.connectors.db2.source;
+
+import org.apache.flink.cdc.common.annotation.Internal;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.connectors.base.config.SourceConfig;
+import org.apache.flink.cdc.connectors.base.source.meta.split.SourceRecords;
+import org.apache.flink.cdc.connectors.base.source.meta.split.SourceSplitState;
+import org.apache.flink.cdc.connectors.base.source.metrics.SourceReaderMetrics;
+import org.apache.flink.cdc.connectors.db2.source.config.Db2SourceConfig;
+import org.apache.flink.cdc.connectors.db2.source.config.Db2SourceConfigFactory;
+import org.apache.flink.cdc.connectors.db2.source.dialect.Db2Dialect;
+import org.apache.flink.cdc.connectors.db2.source.offset.LsnFactory;
+import org.apache.flink.cdc.connectors.db2.source.reader.Db2PipelineRecordEmitter;
+import org.apache.flink.cdc.debezium.DebeziumDeserializationSchema;
+import org.apache.flink.connector.base.source.reader.RecordEmitter;
+
+/**
+ * The Db2 CDC Source for Pipeline connector, which supports parallel reading snapshot of table and
+ * then continue to capture data change from the change log.
+ *
+ * This source extends {@link Db2SourceBuilder.Db2IncrementalSource} and overrides the record
+ * emitter to use {@link Db2PipelineRecordEmitter} for proper handling of schema events in the CDC
+ * pipeline.
+ */
+@Internal
+public class Db2PipelineSource extends Db2SourceBuilder.Db2IncrementalSource {
+
+ private static final long serialVersionUID = 1L;
+
+ public Db2PipelineSource(
+ Db2SourceConfigFactory configFactory,
+ DebeziumDeserializationSchema deserializationSchema,
+ LsnFactory offsetFactory,
+ Db2Dialect dataSourceDialect) {
+ super(configFactory, deserializationSchema, offsetFactory, dataSourceDialect);
+ }
+
+ @Override
+ protected RecordEmitter createRecordEmitter(
+ SourceConfig sourceConfig, SourceReaderMetrics sourceReaderMetrics) {
+ Db2SourceConfig db2SourceConfig = (Db2SourceConfig) sourceConfig;
+ return new Db2PipelineRecordEmitter<>(
+ deserializationSchema, sourceReaderMetrics, db2SourceConfig, offsetFactory);
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SchemaDataTypeInference.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SchemaDataTypeInference.java
new file mode 100644
index 00000000000..846d9901a85
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/Db2SchemaDataTypeInference.java
@@ -0,0 +1,34 @@
+/*
+ * 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.flink.cdc.connectors.db2.source;
+
+import org.apache.flink.cdc.common.annotation.Internal;
+import org.apache.flink.cdc.common.types.DataType;
+import org.apache.flink.cdc.debezium.event.DebeziumSchemaDataTypeInference;
+
+/** {@link DataType} inference for Db2 debezium {@link org.apache.kafka.connect.data.Schema}. */
+@Internal
+public class Db2SchemaDataTypeInference extends DebeziumSchemaDataTypeInference {
+
+ private static final long serialVersionUID = 1L;
+
+ // Db2 has database-specific types, but no special handling is currently
+ // needed here, so this class uses the default implementation from the parent class.
+ // If Db2-specific types require special handling in the future,
+ // it can be added here by overriding the inferStruct method.
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/reader/Db2PipelineRecordEmitter.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/reader/Db2PipelineRecordEmitter.java
new file mode 100644
index 00000000000..02389a0d799
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/source/reader/Db2PipelineRecordEmitter.java
@@ -0,0 +1,203 @@
+/*
+ * 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.flink.cdc.connectors.db2.source.reader;
+
+import org.apache.flink.api.connector.source.SourceOutput;
+import org.apache.flink.cdc.common.event.CreateTableEvent;
+import org.apache.flink.cdc.common.schema.Schema;
+import org.apache.flink.cdc.connectors.base.source.meta.offset.OffsetFactory;
+import org.apache.flink.cdc.connectors.base.source.meta.split.SourceSplitBase;
+import org.apache.flink.cdc.connectors.base.source.meta.split.SourceSplitState;
+import org.apache.flink.cdc.connectors.base.source.metrics.SourceReaderMetrics;
+import org.apache.flink.cdc.connectors.base.source.reader.IncrementalSourceRecordEmitter;
+import org.apache.flink.cdc.connectors.db2.source.config.Db2SourceConfig;
+import org.apache.flink.cdc.connectors.db2.source.dialect.Db2Dialect;
+import org.apache.flink.cdc.connectors.db2.utils.Db2SchemaUtils;
+import org.apache.flink.cdc.debezium.DebeziumDeserializationSchema;
+import org.apache.flink.cdc.debezium.event.DebeziumEventDeserializationSchema;
+import org.apache.flink.connector.base.source.reader.RecordEmitter;
+
+import io.debezium.jdbc.JdbcConnection;
+import io.debezium.relational.history.TableChanges.TableChange;
+import org.apache.kafka.connect.source.SourceRecord;
+
+import java.sql.SQLException;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+
+import static org.apache.flink.cdc.connectors.base.source.meta.wartermark.WatermarkEvent.isLowWatermarkEvent;
+import static org.apache.flink.cdc.connectors.base.utils.SourceRecordUtils.getTableId;
+import static org.apache.flink.cdc.connectors.base.utils.SourceRecordUtils.isDataChangeRecord;
+import static org.apache.flink.cdc.connectors.base.utils.SourceRecordUtils.isSchemaChangeEvent;
+
+/** The {@link RecordEmitter} implementation for Db2 pipeline connector. */
+public class Db2PipelineRecordEmitter extends IncrementalSourceRecordEmitter {
+
+ private final Db2SourceConfig sourceConfig;
+ private final Db2Dialect dataSourceDialect;
+
+ // Track tables that have already sent CreateTableEvent
+ private final Set alreadySendCreateTableTables;
+
+ // Cache for CreateTableEvent, using Map for O(1) lookup
+ private final Map createTableEventCache;
+
+ public Db2PipelineRecordEmitter(
+ DebeziumDeserializationSchema debeziumDeserializationSchema,
+ SourceReaderMetrics sourceReaderMetrics,
+ Db2SourceConfig sourceConfig,
+ OffsetFactory offsetFactory) {
+ super(
+ debeziumDeserializationSchema,
+ sourceReaderMetrics,
+ sourceConfig.isIncludeSchemaChanges(),
+ offsetFactory);
+ this.sourceConfig = sourceConfig;
+ this.dataSourceDialect = new Db2Dialect(sourceConfig);
+ this.alreadySendCreateTableTables = new HashSet<>();
+ this.createTableEventCache =
+ ((DebeziumEventDeserializationSchema) debeziumDeserializationSchema)
+ .getCreateTableEventCache();
+ }
+
+ @Override
+ protected void processElement(
+ SourceRecord element, SourceOutput output, SourceSplitState splitState)
+ throws Exception {
+ if (isSchemaChangeEvent(element) && splitState.isStreamSplitState()) {
+ cacheCreateTableEventsFromSchemas(splitState.asStreamSplitState().getTableSchemas());
+ }
+
+ if (isLowWatermarkEvent(element) && splitState.isSnapshotSplitState()) {
+ // In Snapshot phase of INITIAL startup mode, lazily send CreateTableEvent
+ // to downstream to avoid checkpoint timeout.
+ io.debezium.relational.TableId tableId =
+ splitState.asSnapshotSplitState().toSourceSplit().getTableId();
+ emitCreateTableEventIfNeeded(tableId, output, splitState);
+ } else if (isDataChangeRecord(element)) {
+ // Handle data change events, schema change events are handled downstream directly
+ io.debezium.relational.TableId tableId = getTableId(element);
+ emitCreateTableEventIfNeeded(tableId, output, splitState);
+ }
+ super.processElement(element, output, splitState);
+ }
+
+ @Override
+ public void applySplit(SourceSplitBase split) {
+ cacheCreateTableEventsFromSchemas(split.getTableSchemas());
+ }
+
+ @SuppressWarnings("unchecked")
+ private void emitCreateTableEventIfNeeded(
+ io.debezium.relational.TableId tableId,
+ SourceOutput output,
+ SourceSplitState splitState) {
+ if (alreadySendCreateTableTables.contains(tableId)) {
+ return;
+ }
+
+ cacheCreateTableEventsFromSchemas(splitState.toSourceSplit().getTableSchemas());
+ output.collect((T) getOrCreateCreateTableEvent(tableId, splitState));
+ alreadySendCreateTableTables.add(tableId);
+ }
+
+ private CreateTableEvent getOrCreateCreateTableEvent(
+ io.debezium.relational.TableId tableId, SourceSplitState splitState) {
+ CreateTableEvent createTableEvent = createTableEventCache.get(tableId);
+ if (createTableEvent == null) {
+ createTableEvent = getCreateTableEventFromSplit(tableId, splitState);
+ }
+ if (createTableEvent == null) {
+ createTableEvent = getCreateTableEventFromDatabase(tableId);
+ }
+ createTableEventCache.put(tableId, createTableEvent);
+ return createTableEvent;
+ }
+
+ private CreateTableEvent getCreateTableEventFromSplit(
+ io.debezium.relational.TableId tableId, SourceSplitState splitState) {
+ Map tableSchemas =
+ splitState.toSourceSplit().getTableSchemas();
+ if (tableSchemas == null || tableSchemas.isEmpty()) {
+ return null;
+ }
+ TableChange tableChange = tableSchemas.get(tableId);
+ if (tableChange == null) {
+ // The catalog carried by a split comes from the Db2 system catalog (e.g. "TESTDB"),
+ // while Debezium reports the configured database name in source records (e.g.
+ // "testdb"). The two may differ in case, so fall back to matching on schema and
+ // table name only.
+ for (Map.Entry entry :
+ tableSchemas.entrySet()) {
+ if (matchesIgnoringCatalog(entry.getKey(), tableId)) {
+ tableChange = entry.getValue();
+ break;
+ }
+ }
+ }
+ if (tableChange == null || tableChange.getTable() == null) {
+ return null;
+ }
+ return buildCreateTableEvent(tableId, Db2SchemaUtils.toSchema(tableChange.getTable()));
+ }
+
+ /** Last resort: read the current table schema from the database itself. */
+ private CreateTableEvent getCreateTableEventFromDatabase(
+ io.debezium.relational.TableId tableId) {
+ org.apache.flink.cdc.common.event.TableId cdcTableId = Db2SchemaUtils.toCdcTableId(tableId);
+ try (JdbcConnection jdbc = dataSourceDialect.openJdbcConnection(sourceConfig)) {
+ return buildCreateTableEvent(
+ tableId, Db2SchemaUtils.getTableSchema(cdcTableId, jdbc, dataSourceDialect));
+ } catch (SQLException e) {
+ throw new RuntimeException(
+ "Cannot fetch table schema from database for " + cdcTableId, e);
+ }
+ }
+
+ private static boolean matchesIgnoringCatalog(
+ io.debezium.relational.TableId left, io.debezium.relational.TableId right) {
+ return Objects.equals(left.schema(), right.schema())
+ && Objects.equals(left.table(), right.table());
+ }
+
+ private void cacheCreateTableEventsFromSchemas(
+ Map tableSchemas) {
+ if (tableSchemas == null || tableSchemas.isEmpty()) {
+ return;
+ }
+ for (Map.Entry entry :
+ tableSchemas.entrySet()) {
+ io.debezium.relational.TableId tableId = entry.getKey();
+ TableChange tableChange = entry.getValue();
+ if (tableId == null || tableChange == null || tableChange.getTable() == null) {
+ continue;
+ }
+ createTableEventCache.put(
+ tableId,
+ buildCreateTableEvent(
+ tableId, Db2SchemaUtils.toSchema(tableChange.getTable())));
+ }
+ }
+
+ private CreateTableEvent buildCreateTableEvent(
+ io.debezium.relational.TableId tableId, Schema schema) {
+ return new CreateTableEvent(Db2SchemaUtils.toCdcTableId(tableId), schema);
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2SchemaUtils.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2SchemaUtils.java
new file mode 100644
index 00000000000..9d77474d89f
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2SchemaUtils.java
@@ -0,0 +1,180 @@
+/*
+ * 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.flink.cdc.connectors.db2.utils;
+
+import org.apache.flink.cdc.common.event.TableId;
+import org.apache.flink.cdc.common.schema.Column;
+import org.apache.flink.cdc.common.schema.Schema;
+import org.apache.flink.cdc.connectors.db2.source.config.Db2SourceConfig;
+import org.apache.flink.cdc.connectors.db2.source.dialect.Db2Dialect;
+
+import io.debezium.jdbc.JdbcConnection;
+import io.debezium.relational.Table;
+import io.debezium.relational.history.TableChanges.TableChange;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import javax.annotation.Nullable;
+
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.stream.Collectors;
+
+/** Utilities for converting from debezium {@link Table} types to {@link Schema}. */
+public class Db2SchemaUtils {
+
+ private static final Logger LOG = LoggerFactory.getLogger(Db2SchemaUtils.class);
+
+ private static final String LIST_SCHEMAS_SQL = "SELECT SCHEMANAME FROM SYSCAT.SCHEMATA";
+
+ private static final String LIST_TABLES_SQL =
+ "SELECT TABSCHEMA, TABNAME FROM SYSCAT.TABLES WHERE TYPE = 'T'";
+
+ /** List all schemas of the database in the given {@link Db2SourceConfig}. */
+ public static List listSchemas(Db2SourceConfig sourceConfig) {
+ try (JdbcConnection jdbc = createDb2Connection(sourceConfig)) {
+ return listSchemas(jdbc);
+ } catch (SQLException e) {
+ throw new RuntimeException("Error to list schemas: " + e.getMessage(), e);
+ }
+ }
+
+ public static List listSchemas(JdbcConnection jdbc) throws SQLException {
+ LOG.info("Read list of available schemas");
+ final List schemaNames = new ArrayList<>();
+ jdbc.query(
+ LIST_SCHEMAS_SQL,
+ rs -> {
+ while (rs.next()) {
+ schemaNames.add(rs.getString(1));
+ }
+ });
+ LOG.info("\t list of available schemas are: {}", schemaNames);
+ return schemaNames;
+ }
+
+ /**
+ * List all tables of the database in the given {@link Db2SourceConfig}.
+ *
+ * @param sourceConfig The source configuration.
+ * @param schemaName The schema to list tables from. If null, list tables from all schemas.
+ * @return The list of {@link TableId}s in the format of "schema.table".
+ */
+ public static List listTables(
+ Db2SourceConfig sourceConfig, @Nullable String schemaName) {
+ try (JdbcConnection jdbc = createDb2Connection(sourceConfig)) {
+ return listTables(jdbc, schemaName);
+ } catch (SQLException e) {
+ throw new RuntimeException("Error to list tables: " + e.getMessage(), e);
+ }
+ }
+
+ public static List listTables(JdbcConnection jdbc, @Nullable String schemaName)
+ throws SQLException {
+ LOG.info("Read list of available tables");
+ final List tableIds = new ArrayList<>();
+ String querySql =
+ schemaName == null
+ ? LIST_TABLES_SQL
+ : LIST_TABLES_SQL + " AND TABSCHEMA = '" + schemaName + "'";
+ jdbc.query(
+ querySql,
+ rs -> {
+ while (rs.next()) {
+ tableIds.add(TableId.tableId(rs.getString(1), rs.getString(2)));
+ }
+ });
+ LOG.info("\t list of available tables are: {}", tableIds);
+ return tableIds;
+ }
+
+ /** Get the {@link Schema} of the given table. */
+ public static Schema getTableSchema(Db2SourceConfig sourceConfig, TableId tableId) {
+ Db2Dialect dialect = new Db2Dialect(sourceConfig);
+ try (JdbcConnection jdbc = createDb2Connection(sourceConfig)) {
+ return getTableSchema(tableId, jdbc, dialect);
+ } catch (SQLException e) {
+ throw new RuntimeException("Error to get table schema: " + e.getMessage(), e);
+ }
+ }
+
+ public static Schema getTableSchema(TableId tableId, JdbcConnection jdbc, Db2Dialect dialect) {
+ try {
+ TableChange tableChange = dialect.queryTableSchema(jdbc, toDbzTableId(tableId, jdbc));
+ if (tableChange == null || tableChange.getTable() == null) {
+ throw new RuntimeException("Cannot find table schema for " + tableId);
+ }
+ return toSchema(tableChange.getTable());
+ } catch (Exception e) {
+ throw new RuntimeException("Failed to get table schema for " + tableId, e);
+ }
+ }
+
+ public static Schema toSchema(Table table) {
+ List columns =
+ table.columns().stream().map(Db2SchemaUtils::toColumn).collect(Collectors.toList());
+
+ return Schema.newBuilder()
+ .setColumns(columns)
+ .primaryKey(table.primaryKeyColumnNames())
+ .comment(table.comment())
+ .build();
+ }
+
+ public static Column toColumn(io.debezium.relational.Column column) {
+ if (column.defaultValueExpression().isPresent()) {
+ return Column.physicalColumn(
+ column.name(),
+ Db2TypeUtils.fromDbzColumn(column),
+ column.comment(),
+ column.defaultValueExpression().get());
+ } else {
+ return Column.physicalColumn(
+ column.name(), Db2TypeUtils.fromDbzColumn(column), column.comment());
+ }
+ }
+
+ /**
+ * Converts a pipeline {@link TableId} to a debezium {@link io.debezium.relational.TableId}. The
+ * pipeline TableId of Db2 is in the format of "schema.table", while the catalog of the debezium
+ * TableId is the database name.
+ */
+ public static io.debezium.relational.TableId toDbzTableId(
+ TableId tableId, String databaseName) {
+ return new io.debezium.relational.TableId(
+ databaseName, tableId.getSchemaName(), tableId.getTableName());
+ }
+
+ private static io.debezium.relational.TableId toDbzTableId(
+ TableId tableId, JdbcConnection jdbc) {
+ String databaseName =
+ ((io.debezium.connector.db2.Db2Connection) jdbc).getRealDatabaseName();
+ return toDbzTableId(tableId, databaseName);
+ }
+
+ /** Converts a debezium {@link io.debezium.relational.TableId} to a pipeline {@link TableId}. */
+ public static TableId toCdcTableId(io.debezium.relational.TableId dbzTableId) {
+ return TableId.tableId(dbzTableId.schema(), dbzTableId.table());
+ }
+
+ public static JdbcConnection createDb2Connection(Db2SourceConfig sourceConfig) {
+ Db2Dialect dialect = new Db2Dialect(sourceConfig);
+ return dialect.openJdbcConnection(sourceConfig);
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtils.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtils.java
new file mode 100644
index 00000000000..f003dce35af
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtils.java
@@ -0,0 +1,123 @@
+/*
+ * 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.flink.cdc.connectors.db2.utils;
+
+import org.apache.flink.cdc.common.types.DataType;
+import org.apache.flink.cdc.common.types.DataTypes;
+
+import io.debezium.relational.Column;
+
+import java.sql.Types;
+
+/** A utility class for converting Db2 types to Flink CDC types. */
+public class Db2TypeUtils {
+
+ // Db2 specific type names
+ static final String DECFLOAT = "decfloat";
+ static final String XML = "xml";
+
+ private static final int DECFLOAT_PRECISION_16 = 16;
+ private static final int DECFLOAT_PRECISION_34 = 34;
+
+ /** Returns a corresponding Flink CDC data type from a debezium {@link Column}. */
+ public static DataType fromDbzColumn(Column column) {
+ DataType dataType = convertFromColumn(column);
+ if (column.isOptional()) {
+ return dataType;
+ } else {
+ return dataType.notNull();
+ }
+ }
+
+ /**
+ * Returns a corresponding Flink CDC data type from a debezium {@link Column} with nullable
+ * always be true.
+ */
+ private static DataType convertFromColumn(Column column) {
+ int precision = column.length();
+ int scale = column.scale().orElse(0);
+ String typeName = column.typeName();
+
+ if (typeName != null) {
+ switch (typeName.toLowerCase()) {
+ case DECFLOAT:
+ // DECFLOAT(16) maps to DOUBLE, DECFLOAT(34) maps to DECIMAL(34, 0)
+ if (precision == DECFLOAT_PRECISION_34) {
+ return DataTypes.DECIMAL(DECFLOAT_PRECISION_34, 0);
+ }
+ return DataTypes.DOUBLE();
+ case XML:
+ return DataTypes.STRING();
+ default:
+ // Fall through to JDBC type handling.
+ }
+ }
+
+ switch (column.jdbcType()) {
+ case Types.CHAR:
+ if (precision > 0) {
+ return DataTypes.CHAR(precision);
+ }
+ return DataTypes.STRING();
+ case Types.VARCHAR:
+ case Types.LONGVARCHAR:
+ if (precision > 0) {
+ return DataTypes.VARCHAR(precision);
+ }
+ return DataTypes.STRING();
+ case Types.SQLXML:
+ case Types.CLOB:
+ return DataTypes.STRING();
+ case Types.BLOB:
+ case Types.BINARY:
+ case Types.VARBINARY:
+ case Types.LONGVARBINARY:
+ return DataTypes.BYTES();
+ case Types.TINYINT:
+ case Types.SMALLINT:
+ // Db2 SMALLINT is a 2-byte integer
+ return DataTypes.SMALLINT();
+ case Types.INTEGER:
+ return DataTypes.INT();
+ case Types.BIGINT:
+ return DataTypes.BIGINT();
+ case Types.REAL:
+ return DataTypes.FLOAT();
+ case Types.FLOAT:
+ case Types.DOUBLE:
+ return DataTypes.DOUBLE();
+ case Types.DECIMAL:
+ case Types.NUMERIC:
+ if (precision > 0) {
+ return DataTypes.DECIMAL(precision, scale);
+ }
+ return DataTypes.DECIMAL(38, scale);
+ case Types.DATE:
+ return DataTypes.DATE();
+ case Types.TIME:
+ return DataTypes.TIME(Math.max(scale, 0));
+ case Types.TIMESTAMP:
+ return DataTypes.TIMESTAMP(column.scale().orElse(6));
+ default:
+ throw new UnsupportedOperationException(
+ String.format(
+ "Doesn't support Db2 type '%s', JDBC type '%d' yet.",
+ column.typeName(), column.jdbcType()));
+ }
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/resources/META-INF/services/org.apache.flink.cdc.common.factories.Factory b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/resources/META-INF/services/org.apache.flink.cdc.common.factories.Factory
new file mode 100644
index 00000000000..a0afa50534a
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/main/resources/META-INF/services/org.apache.flink.cdc.common.factories.Factory
@@ -0,0 +1,16 @@
+# 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.
+
+org.apache.flink.cdc.connectors.db2.factory.Db2DataSourceFactory
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/factory/Db2DataSourceFactoryTest.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/factory/Db2DataSourceFactoryTest.java
new file mode 100644
index 00000000000..82c1d1c7876
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/factory/Db2DataSourceFactoryTest.java
@@ -0,0 +1,248 @@
+/*
+ * 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.flink.cdc.connectors.db2.factory;
+
+import org.apache.flink.cdc.common.configuration.ConfigOption;
+import org.apache.flink.cdc.common.configuration.Configuration;
+import org.apache.flink.cdc.common.factories.Factory;
+import org.apache.flink.cdc.connectors.db2.Db2TestBase;
+import org.apache.flink.cdc.connectors.db2.source.Db2DataSource;
+import org.apache.flink.table.api.ValidationException;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.Db2Container;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.DATABASE;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.HOSTNAME;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.METADATA_LIST;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.PASSWORD;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.PORT;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.SCAN_INCREMENTAL_SNAPSHOT_CHUNK_KEY_COLUMN;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.SCAN_STARTUP_MODE;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.TABLES;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.TABLES_EXCLUDE;
+import static org.apache.flink.cdc.connectors.db2.source.Db2DataSourceOptions.USERNAME;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests for {@link Db2DataSourceFactory}. */
+public class Db2DataSourceFactoryTest extends Db2TestBase {
+
+ private static final String SCHEMA_NAME = "DB2INST1";
+
+ private static final String TABLE_NAME = "CUSTOMERS";
+
+ @BeforeEach
+ public void before() {
+ initializeDb2Table("customers", TABLE_NAME);
+ }
+
+ private Map getBaseOptions() {
+ Map options = new HashMap<>();
+ options.put(HOSTNAME.key(), DB2_CONTAINER.getHost());
+ options.put(PORT.key(), String.valueOf(DB2_CONTAINER.getMappedPort(Db2Container.DB2_PORT)));
+ options.put(USERNAME.key(), DB2_CONTAINER.getUsername());
+ options.put(PASSWORD.key(), DB2_CONTAINER.getPassword());
+ options.put(DATABASE.key(), DB2_CONTAINER.getDatabaseName());
+ options.put(TABLES.key(), SCHEMA_NAME + "." + TABLE_NAME);
+ return options;
+ }
+
+ @Test
+ public void testCreateDataSource() {
+ Factory.Context context = new MockContext(Configuration.fromMap(getBaseOptions()));
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Db2DataSource dataSource = (Db2DataSource) factory.createDataSource(context);
+ assertThat(dataSource.getDb2SourceConfig().getTableList())
+ .containsExactly(SCHEMA_NAME + "." + TABLE_NAME);
+ }
+
+ @Test
+ public void testNoMatchedTable() {
+ Map options = getBaseOptions();
+ String tables = SCHEMA_NAME + ".nonexistent";
+ options.put(TABLES.key(), tables);
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Cannot find any table by the option 'tables' = " + tables);
+ }
+
+ @Test
+ public void testExcludeAllTable() {
+ Map options = getBaseOptions();
+ String tablesExclude = SCHEMA_NAME + "." + TABLE_NAME;
+ options.put(TABLES_EXCLUDE.key(), tablesExclude);
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining(
+ "Cannot find any table with by the option 'tables.exclude' = "
+ + tablesExclude);
+ }
+
+ @Test
+ public void testLackRequireOption() {
+ Map options = getBaseOptions();
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ List requireKeys =
+ factory.requiredOptions().stream()
+ .map(ConfigOption::key)
+ .collect(Collectors.toList());
+ for (String requireKey : requireKeys) {
+ Map remainingOptions = new HashMap<>(options);
+ remainingOptions.remove(requireKey);
+ Factory.Context context = new MockContext(Configuration.fromMap(remainingOptions));
+
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(ValidationException.class)
+ .hasMessageContaining(
+ String.format(
+ "One or more required options are missing.\n\n"
+ + "Missing required options are:\n\n"
+ + "%s",
+ requireKey));
+ }
+ }
+
+ @Test
+ public void testUnsupportedOption() {
+ Map options = getBaseOptions();
+ options.put("unsupported_key", "unsupported_value");
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(ValidationException.class)
+ .hasMessageContaining(
+ "Unsupported options found for 'db2'.\n\n"
+ + "Unsupported options:\n\n"
+ + "unsupported_key");
+ }
+
+ @Test
+ public void testOptionalOption() {
+ Map options = getBaseOptions();
+
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ assertThat(factory.optionalOptions()).contains(PORT);
+
+ Db2DataSource dataSource = (Db2DataSource) factory.createDataSource(context);
+ assertThat(dataSource.getDb2SourceConfig().getPort())
+ .isEqualTo(DB2_CONTAINER.getMappedPort(Db2Container.DB2_PORT));
+ }
+
+ @Test
+ public void testChunkKeyColumnOptionIsSupported() {
+ Map options = getBaseOptions();
+ options.put(SCAN_INCREMENTAL_SNAPSHOT_CHUNK_KEY_COLUMN.key(), "ID");
+
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+
+ assertThat(factory.optionalOptions()).contains(SCAN_INCREMENTAL_SNAPSHOT_CHUNK_KEY_COLUMN);
+ Db2DataSource dataSource = (Db2DataSource) factory.createDataSource(context);
+ assertThat(dataSource.getDb2SourceConfig().getChunkKeyColumn()).isEqualTo("ID");
+ }
+
+ @Test
+ public void testUnsupportedStartupMode() {
+ Map options = getBaseOptions();
+ options.put(SCAN_STARTUP_MODE.key(), "timestamp");
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(ValidationException.class)
+ .hasMessageContaining("Invalid value for option 'scan.startup.mode'");
+ }
+
+ @Test
+ public void testLatestOffsetStartupMode() {
+ Map options = getBaseOptions();
+ options.put(SCAN_STARTUP_MODE.key(), "latest-offset");
+
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Db2DataSource dataSource = (Db2DataSource) factory.createDataSource(context);
+ assertThat(dataSource.getDb2SourceConfig().getStartupOptions().isStreamOnly()).isTrue();
+ }
+
+ @Test
+ public void testPrefixRequireOption() {
+ Map options = getBaseOptions();
+ options.put("debezium.snapshot.mode", "initial");
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Db2DataSource dataSource = (Db2DataSource) factory.createDataSource(context);
+ assertThat(dataSource.getDb2SourceConfig().getTableList())
+ .containsExactly(SCHEMA_NAME + "." + TABLE_NAME);
+ }
+
+ @Test
+ public void testInvalidMetadataList() {
+ Map options = getBaseOptions();
+ options.put(METADATA_LIST.key(), "database_name,unknown_metadata");
+
+ Db2DataSourceFactory factory = new Db2DataSourceFactory();
+ Factory.Context context = new MockContext(Configuration.fromMap(options));
+
+ assertThatThrownBy(() -> factory.createDataSource(context))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("cannot be found in Db2 metadata");
+ }
+
+ static class MockContext implements Factory.Context {
+
+ Configuration factoryConfiguration;
+
+ public MockContext(Configuration factoryConfiguration) {
+ this.factoryConfiguration = factoryConfiguration;
+ }
+
+ @Override
+ public Configuration getFactoryConfiguration() {
+ return factoryConfiguration;
+ }
+
+ @Override
+ public Configuration getPipelineConfiguration() {
+ return null;
+ }
+
+ @Override
+ public ClassLoader getClassLoader() {
+ return this.getClass().getClassLoader();
+ }
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineITCase.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineITCase.java
new file mode 100644
index 00000000000..8e1d9aa72cf
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/source/Db2PipelineITCase.java
@@ -0,0 +1,314 @@
+/*
+ * 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.flink.cdc.connectors.db2.source;
+
+import org.apache.flink.api.common.eventtime.WatermarkStrategy;
+import org.apache.flink.cdc.common.data.binary.BinaryStringData;
+import org.apache.flink.cdc.common.event.CreateTableEvent;
+import org.apache.flink.cdc.common.event.DataChangeEvent;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.common.event.TableId;
+import org.apache.flink.cdc.common.factories.Factory;
+import org.apache.flink.cdc.common.factories.FactoryHelper;
+import org.apache.flink.cdc.common.schema.Schema;
+import org.apache.flink.cdc.common.source.FlinkSourceProvider;
+import org.apache.flink.cdc.common.types.DataType;
+import org.apache.flink.cdc.common.types.DataTypes;
+import org.apache.flink.cdc.common.types.RowType;
+import org.apache.flink.cdc.connectors.base.options.StartupOptions;
+import org.apache.flink.cdc.connectors.db2.Db2TestBase;
+import org.apache.flink.cdc.connectors.db2.factory.Db2DataSourceFactory;
+import org.apache.flink.cdc.connectors.db2.source.config.Db2SourceConfigFactory;
+import org.apache.flink.cdc.runtime.typeutils.BinaryRecordDataGenerator;
+import org.apache.flink.cdc.runtime.typeutils.EventTypeInfo;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.util.RestartStrategyUtils;
+import org.apache.flink.util.CloseableIterator;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.Db2Container;
+
+import java.sql.Connection;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Integration tests for Db2 pipeline source. */
+public class Db2PipelineITCase extends Db2TestBase {
+
+ private static final String SCHEMA_NAME = "DB2INST1";
+
+ private static final String TABLE_NAME = "CUSTOMERS";
+
+ private static final TableId TABLE_ID = TableId.tableId(SCHEMA_NAME, TABLE_NAME);
+
+ private static final StreamExecutionEnvironment env =
+ StreamExecutionEnvironment.getExecutionEnvironment();
+
+ @SuppressWarnings("deprecation")
+ @BeforeEach
+ public void before() {
+ env.setParallelism(4);
+ env.enableCheckpointing(200);
+ RestartStrategyUtils.configureNoRestartStrategy(env);
+ initializeDb2Table("customers", TABLE_NAME);
+ }
+
+ @Test
+ public void testInitialStartupMode() throws Exception {
+ Db2SourceConfigFactory configFactory =
+ (Db2SourceConfigFactory)
+ new Db2SourceConfigFactory()
+ .hostname(DB2_CONTAINER.getHost())
+ .port(DB2_CONTAINER.getMappedPort(Db2Container.DB2_PORT))
+ .username(DB2_CONTAINER.getUsername())
+ .password(DB2_CONTAINER.getPassword())
+ .databaseList(DB2_CONTAINER.getDatabaseName())
+ .tableList(SCHEMA_NAME + "." + TABLE_NAME)
+ .startupOptions(StartupOptions.initial())
+ .serverTimeZone("UTC");
+
+ FlinkSourceProvider sourceProvider =
+ (FlinkSourceProvider) new Db2DataSource(configFactory).getEventSourceProvider();
+ CloseableIterator events =
+ env.fromSource(
+ sourceProvider.getSource(),
+ WatermarkStrategy.noWatermarks(),
+ Db2DataSourceFactory.IDENTIFIER,
+ new EventTypeInfo())
+ .executeAndCollect();
+
+ CreateTableEvent createTableEvent = getCustomersCreateTableEvent();
+ List expectedSnapshot = getSnapshotExpected();
+
+ List actual = fetchResultsExcept(events, expectedSnapshot.size(), createTableEvent);
+ assertThat(actual.subList(0, expectedSnapshot.size()))
+ .containsExactlyInAnyOrder(expectedSnapshot.toArray(new Event[0]));
+ }
+
+ @Test
+ public void testInitialStartupModeWithMetadata() throws Exception {
+ org.apache.flink.cdc.common.configuration.Configuration sourceConfiguration =
+ new org.apache.flink.cdc.common.configuration.Configuration();
+ sourceConfiguration.set(Db2DataSourceOptions.HOSTNAME, DB2_CONTAINER.getHost());
+ sourceConfiguration.set(
+ Db2DataSourceOptions.PORT, DB2_CONTAINER.getMappedPort(Db2Container.DB2_PORT));
+ sourceConfiguration.set(Db2DataSourceOptions.USERNAME, DB2_CONTAINER.getUsername());
+ sourceConfiguration.set(Db2DataSourceOptions.PASSWORD, DB2_CONTAINER.getPassword());
+ sourceConfiguration.set(Db2DataSourceOptions.DATABASE, DB2_CONTAINER.getDatabaseName());
+ sourceConfiguration.set(Db2DataSourceOptions.TABLES, SCHEMA_NAME + "." + TABLE_NAME);
+ sourceConfiguration.set(Db2DataSourceOptions.SERVER_TIME_ZONE, "UTC");
+ sourceConfiguration.set(
+ Db2DataSourceOptions.METADATA_LIST, "database_name,schema_name,table_name,op_ts");
+
+ Factory.Context context =
+ new FactoryHelper.DefaultContext(
+ sourceConfiguration,
+ new org.apache.flink.cdc.common.configuration.Configuration(),
+ this.getClass().getClassLoader());
+ FlinkSourceProvider sourceProvider =
+ (FlinkSourceProvider)
+ new Db2DataSourceFactory()
+ .createDataSource(context)
+ .getEventSourceProvider();
+ CloseableIterator events =
+ env.fromSource(
+ sourceProvider.getSource(),
+ WatermarkStrategy.noWatermarks(),
+ Db2DataSourceFactory.IDENTIFIER,
+ new EventTypeInfo())
+ .executeAndCollect();
+
+ CreateTableEvent createTableEvent = getCustomersCreateTableEvent();
+
+ Map meta = new HashMap<>();
+ meta.put("database_name", DB2_CONTAINER.getDatabaseName());
+ meta.put("schema_name", SCHEMA_NAME);
+ meta.put("table_name", TABLE_NAME);
+ meta.put("op_ts", "0");
+
+ List expectedSnapshot =
+ getSnapshotExpected().stream()
+ .map(
+ event -> {
+ DataChangeEvent dataChangeEvent = (DataChangeEvent) event;
+ return DataChangeEvent.insertEvent(
+ dataChangeEvent.tableId(),
+ dataChangeEvent.after(),
+ meta);
+ })
+ .collect(Collectors.toList());
+
+ List expectedLog = new ArrayList<>();
+
+ RowType rowType = getCustomersRowType();
+ BinaryRecordDataGenerator generator = new BinaryRecordDataGenerator(rowType);
+
+ try (Connection connection = getJdbcConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute(
+ "INSERT INTO "
+ + SCHEMA_NAME
+ + "."
+ + TABLE_NAME
+ + " VALUES (1031, 'user_22', 'Berlin', '123567891234')");
+ expectedLog.add(
+ DataChangeEvent.insertEvent(
+ TABLE_ID,
+ generator.generate(
+ new Object[] {
+ 1031,
+ BinaryStringData.fromString("user_22"),
+ BinaryStringData.fromString("Berlin"),
+ BinaryStringData.fromString("123567891234")
+ })));
+ statement.execute(
+ "UPDATE "
+ + SCHEMA_NAME
+ + "."
+ + TABLE_NAME
+ + " SET ADDRESS='Hangzhou' WHERE ID = 1031");
+ expectedLog.add(
+ DataChangeEvent.updateEvent(
+ TABLE_ID,
+ generator.generate(
+ new Object[] {
+ 1031,
+ BinaryStringData.fromString("user_22"),
+ BinaryStringData.fromString("Berlin"),
+ BinaryStringData.fromString("123567891234")
+ }),
+ generator.generate(
+ new Object[] {
+ 1031,
+ BinaryStringData.fromString("user_22"),
+ BinaryStringData.fromString("Hangzhou"),
+ BinaryStringData.fromString("123567891234")
+ })));
+ statement.execute("DELETE FROM " + SCHEMA_NAME + "." + TABLE_NAME + " WHERE ID = 1031");
+ expectedLog.add(
+ DataChangeEvent.deleteEvent(
+ TABLE_ID,
+ generator.generate(
+ new Object[] {
+ 1031,
+ BinaryStringData.fromString("user_22"),
+ BinaryStringData.fromString("Hangzhou"),
+ BinaryStringData.fromString("123567891234")
+ })));
+ }
+
+ int snapshotRecordsCount = expectedSnapshot.size();
+ int logRecordsCount = expectedLog.size();
+
+ List actual =
+ fetchResultsExcept(
+ events, snapshotRecordsCount + logRecordsCount, createTableEvent);
+
+ List actualSnapshotEvents = actual.subList(0, snapshotRecordsCount);
+ List actualLogEvents = actual.subList(snapshotRecordsCount, actual.size());
+
+ assertThat(actualSnapshotEvents).containsExactlyInAnyOrderElementsOf(expectedSnapshot);
+ assertThat(actualLogEvents).hasSize(logRecordsCount);
+
+ for (int i = 0; i < logRecordsCount; i++) {
+ DataChangeEvent expectedEvent = (DataChangeEvent) expectedLog.get(i);
+ DataChangeEvent actualEvent = (DataChangeEvent) actualLogEvents.get(i);
+ assertThat(actualEvent.op()).isEqualTo(expectedEvent.op());
+ assertThat(actualEvent.before()).isEqualTo(expectedEvent.before());
+ assertThat(actualEvent.after()).isEqualTo(expectedEvent.after());
+ assertThat(actualEvent.meta().get("database_name"))
+ .isEqualTo(DB2_CONTAINER.getDatabaseName());
+ assertThat(actualEvent.meta().get("schema_name")).isEqualTo(SCHEMA_NAME);
+ assertThat(actualEvent.meta().get("table_name")).isEqualTo(TABLE_NAME);
+ // The Db2 connector resolves LSN timestamps with the JVM default time zone,
+ // which CI randomizes, so only assert op_ts is a valid change timestamp here.
+ assertThat(Long.parseLong(actualEvent.meta().get("op_ts"))).isGreaterThan(0L);
+ }
+ }
+
+ private static List fetchResultsExcept(Iterator iter, int size, T sideEvent) {
+ List result = new ArrayList<>(size);
+ List sideResults = new ArrayList<>();
+ while (size > 0 && iter.hasNext()) {
+ T event = iter.next();
+ if (sideEvent.getClass().isInstance(event)) {
+ sideResults.add(event);
+ } else {
+ result.add(event);
+ size--;
+ }
+ }
+ // Also ensure we've received at least one or many side events.
+ assertThat(sideResults).isNotEmpty();
+ return result;
+ }
+
+ private static RowType getCustomersRowType() {
+ return RowType.of(
+ new DataType[] {
+ DataTypes.INT().notNull(),
+ DataTypes.VARCHAR(255).notNull(),
+ DataTypes.VARCHAR(1024),
+ DataTypes.VARCHAR(512)
+ },
+ new String[] {"ID", "NAME", "ADDRESS", "PHONE_NUMBER"});
+ }
+
+ private List getSnapshotExpected() {
+ BinaryRecordDataGenerator generator = new BinaryRecordDataGenerator(getCustomersRowType());
+ List snapshotExpected = new ArrayList<>();
+ int[] ids = {
+ 101, 102, 103, 109, 110, 111, 118, 121, 123, 1009, 1010, 1011, 1012, 1013, 1014, 1015,
+ 1016, 1017, 1018, 1019, 2000
+ };
+ for (int i = 0; i < ids.length; i++) {
+ snapshotExpected.add(
+ DataChangeEvent.insertEvent(
+ TABLE_ID,
+ generator.generate(
+ new Object[] {
+ ids[i],
+ BinaryStringData.fromString("user_" + (i + 1)),
+ BinaryStringData.fromString("Shanghai"),
+ BinaryStringData.fromString("123567891234")
+ })));
+ }
+ return snapshotExpected;
+ }
+
+ private CreateTableEvent getCustomersCreateTableEvent() {
+ return new CreateTableEvent(
+ TABLE_ID,
+ Schema.newBuilder()
+ .physicalColumn("ID", DataTypes.INT().notNull())
+ .physicalColumn("NAME", DataTypes.VARCHAR(255).notNull())
+ .physicalColumn("ADDRESS", DataTypes.VARCHAR(1024))
+ .physicalColumn("PHONE_NUMBER", DataTypes.VARCHAR(512))
+ .primaryKey(Collections.singletonList("ID"))
+ .build());
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtilsTest.java b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtilsTest.java
new file mode 100644
index 00000000000..a34eadf7f59
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/java/org/apache/flink/cdc/connectors/db2/utils/Db2TypeUtilsTest.java
@@ -0,0 +1,143 @@
+/*
+ * 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.flink.cdc.connectors.db2.utils;
+
+import org.apache.flink.cdc.common.types.DataTypes;
+
+import io.debezium.relational.Column;
+import io.debezium.relational.ColumnEditor;
+import org.junit.jupiter.api.Test;
+
+import java.sql.Types;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+/** Tests for {@link Db2TypeUtils}. */
+class Db2TypeUtilsTest {
+
+ @Test
+ void testDecFloatMapsToDoubleOrDecimal() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.OTHER, "DECFLOAT", 16, null)))
+ .isEqualTo(DataTypes.DOUBLE());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.OTHER, "DECFLOAT", 34, null)))
+ .isEqualTo(DataTypes.DECIMAL(34, 0));
+ }
+
+ @Test
+ void testXmlMapsToString() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.SQLXML, "XML", 0, null)))
+ .isEqualTo(DataTypes.STRING());
+ }
+
+ @Test
+ void testCharAndVarcharPreserveLength() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.CHAR, "CHAR", 10, null)))
+ .isEqualTo(DataTypes.CHAR(10));
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.VARCHAR, "VARCHAR", 128, null)))
+ .isEqualTo(DataTypes.VARCHAR(128));
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.CHAR, "CHAR", 0, null)))
+ .isEqualTo(DataTypes.STRING());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.VARCHAR, "VARCHAR", 0, null)))
+ .isEqualTo(DataTypes.STRING());
+ }
+
+ @Test
+ void testNumericTypes() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.SMALLINT, "SMALLINT", 0, null)))
+ .isEqualTo(DataTypes.SMALLINT());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.INTEGER, "INTEGER", 0, null)))
+ .isEqualTo(DataTypes.INT());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.BIGINT, "BIGINT", 0, null)))
+ .isEqualTo(DataTypes.BIGINT());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.REAL, "REAL", 0, null)))
+ .isEqualTo(DataTypes.FLOAT());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.DOUBLE, "DOUBLE", 0, null)))
+ .isEqualTo(DataTypes.DOUBLE());
+ }
+
+ @Test
+ void testDecimalPreservesPrecisionAndScale() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.DECIMAL, "DECIMAL", 10, 2)))
+ .isEqualTo(DataTypes.DECIMAL(10, 2));
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.NUMERIC, "NUMERIC", 31, 0)))
+ .isEqualTo(DataTypes.DECIMAL(31, 0));
+ }
+
+ @Test
+ void testTemporalTypes() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.DATE, "DATE", 0, null)))
+ .isEqualTo(DataTypes.DATE());
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.TIME, "TIME", 0, null)))
+ .isEqualTo(DataTypes.TIME(0));
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.TIMESTAMP, "TIMESTAMP", 0, 6)))
+ .isEqualTo(DataTypes.TIMESTAMP(6));
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.TIMESTAMP, "TIMESTAMP", 0, null)))
+ .isEqualTo(DataTypes.TIMESTAMP(6));
+ }
+
+ @Test
+ void testBinaryTypes() {
+ assertThat(Db2TypeUtils.fromDbzColumn(column(Types.BLOB, "BLOB", 0, null)))
+ .isEqualTo(DataTypes.BYTES());
+ assertThat(
+ Db2TypeUtils.fromDbzColumn(
+ column(Types.BINARY, "CHAR () FOR BIT DATA", 8, null)))
+ .isEqualTo(DataTypes.BYTES());
+ assertThat(
+ Db2TypeUtils.fromDbzColumn(
+ column(Types.VARBINARY, "VARCHAR () FOR BIT DATA", 8, null)))
+ .isEqualTo(DataTypes.BYTES());
+ }
+
+ @Test
+ void testNotNullableColumn() {
+ Column notNullColumn =
+ Column.editor()
+ .name("c")
+ .jdbcType(Types.INTEGER)
+ .type("INTEGER")
+ .optional(false)
+ .create();
+ assertThat(Db2TypeUtils.fromDbzColumn(notNullColumn)).isEqualTo(DataTypes.INT().notNull());
+ }
+
+ @Test
+ void testUnsupportedType() {
+ assertThatThrownBy(
+ () ->
+ Db2TypeUtils.fromDbzColumn(
+ column(Types.JAVA_OBJECT, "unknown_type", 0, null)))
+ .isInstanceOf(UnsupportedOperationException.class)
+ .hasMessageContaining("unknown_type");
+ }
+
+ private static Column column(int jdbcType, String typeName, int length, Integer scale) {
+ ColumnEditor editor =
+ Column.editor()
+ .name("c")
+ .jdbcType(jdbcType)
+ .type(typeName)
+ .length(length)
+ .optional(true);
+ if (scale != null) {
+ editor.scale(scale);
+ }
+ return editor.create();
+ }
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/.gitattributes b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/.gitattributes
new file mode 100644
index 00000000000..90716bd2e3e
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/.gitattributes
@@ -0,0 +1,3 @@
+# Files in this directory are copied into a Linux container and executed there,
+# so they must always keep LF line endings, even on Windows checkouts.
+* text=auto eol=lf
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/Dockerfile b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/Dockerfile
new file mode 100644
index 00000000000..fd8db80d30a
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/Dockerfile
@@ -0,0 +1,36 @@
+################################################################################
+# 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.
+################################################################################
+FROM ibmcom/db2:11.5.0.0a
+
+MAINTAINER Peter Urbanetz
+
+RUN mkdir -p /asncdctools/src
+
+ADD asncdc_UDF.sql /asncdctools/src
+ADD asncdcaddremove.sql /asncdctools/src
+ADD asncdctables.sql /asncdctools/src
+ADD dbsetup.sh /asncdctools/src
+ADD startup-agent.sql /asncdctools/src
+ADD inventory.sql /asncdctools/src
+ADD column_type_test.sql /asncdctools/src
+ADD asncdc.c /asncdctools/src
+
+RUN chmod -R 777 /asncdctools
+
+RUN mkdir /var/custom
+ADD cdcsetup.sh /var/custom
+RUN chmod -R 777 /var/custom
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc.c b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc.c
new file mode 100644
index 00000000000..78c5515455d
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc.c
@@ -0,0 +1,179 @@
+/*
+ * 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.
+ */
+
+#include
+#include
+#include
+#include
+#include
+#include
+
+void SQL_API_FN asncdcservice(
+ SQLUDF_VARCHAR *asnCommand, /* input */
+ SQLUDF_VARCHAR *asnService,
+ SQLUDF_CLOB *fileData, /* output */
+ /* null indicators */
+ SQLUDF_NULLIND *asnCommand_ind, /* input */
+ SQLUDF_NULLIND *asnService_ind,
+ SQLUDF_NULLIND *fileData_ind,
+ SQLUDF_TRAIL_ARGS,
+ struct sqludf_dbinfo *dbinfo)
+{
+
+ int fd;
+ char tmpFileName[] = "/tmp/fileXXXXXX";
+ fd = mkstemp(tmpFileName);
+
+ int strcheck = 0;
+ char cmdstring[256];
+
+
+ char* szDb2path = getenv("HOME");
+
+
+
+ char str[20];
+ int len = 0;
+ char c;
+ char *buffer = NULL;
+ FILE *pidfile;
+
+ char dbname[129];
+ memset(dbname, '\0', 129);
+ strncpy(dbname, (char *)(dbinfo->dbname), dbinfo->dbnamelen);
+ dbname[dbinfo->dbnamelen] = '\0';
+
+ int pid;
+ if (strcmp(asnService, "asncdc") == 0)
+ {
+ strcheck = sprintf(cmdstring, "pgrep -fx \"%s/sqllib/bin/asncap capture_schema=%s capture_server=%s\" > %s", szDb2path, asnService, dbname, tmpFileName);
+ int callcheck;
+ callcheck = system(cmdstring);
+ pidfile = fopen(tmpFileName, "r");
+ while ((c = fgetc(pidfile)) != EOF)
+ {
+ if (c == '\n')
+ {
+ break;
+ }
+ len++;
+ }
+ buffer = (char *)malloc(sizeof(char) * len);
+ fseek(pidfile, 0, SEEK_SET);
+ fread(buffer, sizeof(char), len, pidfile);
+ fclose(pidfile);
+ pidfile = fopen(tmpFileName, "w");
+ if (strcmp(asnCommand, "start") == 0)
+ {
+ if (len == 0) // is not running
+ {
+ strcheck = sprintf(cmdstring, "%s/sqllib/bin/asncap capture_schema=%s capture_server=%s &", szDb2path, asnService, dbname);
+ fprintf(pidfile, "start --> %s \n", cmdstring);
+ callcheck = system(cmdstring);
+ }
+ else
+ {
+ fprintf(pidfile, "asncap is already running");
+ }
+ }
+ if ((strcmp(asnCommand, "prune") == 0) ||
+ (strcmp(asnCommand, "reinit") == 0) ||
+ (strcmp(asnCommand, "suspend") == 0) ||
+ (strcmp(asnCommand, "resume") == 0) ||
+ (strcmp(asnCommand, "status") == 0) ||
+ (strcmp(asnCommand, "stop") == 0))
+ {
+ if (len > 0)
+ {
+ //buffer[len] = '\0';
+ //strcheck = sprintf(cmdstring, "/bin/kill -SIGINT %s ", buffer);
+ //fprintf(pidfile, "stop --> %s", cmdstring);
+ //callcheck = system(cmdstring);
+ strcheck = sprintf(cmdstring, "%s/sqllib/bin/asnccmd capture_schema=%s capture_server=%s %s >> %s", szDb2path, asnService, dbname, asnCommand, tmpFileName);
+ //fprintf(pidfile, "%s --> %s \n", cmdstring, asnCommand);
+ callcheck = system(cmdstring);
+ }
+ else
+ {
+ fprintf(pidfile, "asncap is not running");
+ }
+ }
+
+ fclose(pidfile);
+ }
+ /* system(cmdstring); */
+
+ int rc = 0;
+ long fileSize = 0;
+ size_t readCnt = 0;
+ FILE *f = NULL;
+
+ f = fopen(tmpFileName, "r");
+ if (!f)
+ {
+ strcpy(SQLUDF_MSGTX, "Could not open file ");
+ strncat(SQLUDF_MSGTX, tmpFileName,
+ SQLUDF_MSGTEXT_LEN - strlen(SQLUDF_MSGTX) - 1);
+ strncpy(SQLUDF_STATE, "38100", SQLUDF_SQLSTATE_LEN);
+ return;
+ }
+
+ rc = fseek(f, 0, SEEK_END);
+ if (rc)
+ {
+ sprintf(SQLUDF_MSGTX, "fseek() failed with rc = %d", rc);
+ strncpy(SQLUDF_STATE, "38101", SQLUDF_SQLSTATE_LEN);
+ return;
+ }
+
+ /* verify the file size */
+ fileSize = ftell(f);
+ if (fileSize > fileData->length)
+ {
+ strcpy(SQLUDF_MSGTX, "File too large");
+ strncpy(SQLUDF_STATE, "38102", SQLUDF_SQLSTATE_LEN);
+ return;
+ }
+
+ /* go to the beginning and read the entire file */
+ rc = fseek(f, 0, 0);
+ if (rc)
+ {
+ sprintf(SQLUDF_MSGTX, "fseek() failed with rc = %d", rc);
+ strncpy(SQLUDF_STATE, "38103", SQLUDF_SQLSTATE_LEN);
+ return;
+ }
+
+ readCnt = fread(fileData->data, 1, fileSize, f);
+ if (readCnt != fileSize)
+ {
+ /* raise a warning that something weird is going on */
+ sprintf(SQLUDF_MSGTX, "Could not read entire file "
+ "(%d vs %d)",
+ readCnt, fileSize);
+ strncpy(SQLUDF_STATE, "01H10", SQLUDF_SQLSTATE_LEN);
+ *fileData_ind = -1;
+ }
+ else
+ {
+ fileData->length = readCnt;
+ *fileData_ind = 0;
+ }
+ // remove temorary file
+ rc = remove(tmpFileName);
+ //fclose(pFile);
+}
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc_UDF.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc_UDF.sql
new file mode 100644
index 00000000000..fd07e00edc4
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdc_UDF.sql
@@ -0,0 +1,32 @@
+-- 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.
+
+DROP SPECIFIC FUNCTION ASNCDC.asncdcservice;
+
+CREATE FUNCTION ASNCDC.ASNCDCSERVICES(command VARCHAR(6), service VARCHAR(8))
+ RETURNS CLOB(100K)
+ SPECIFIC asncdcservice
+ EXTERNAL NAME 'asncdc!asncdcservice'
+ LANGUAGE C
+ PARAMETER STYLE SQL
+ DBINFO
+ DETERMINISTIC
+ NOT FENCED
+ RETURNS NULL ON NULL INPUT
+ NO SQL
+ NO EXTERNAL ACTION
+ NO SCRATCHPAD
+ ALLOW PARALLEL
+ NO FINAL CALL;
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdcaddremove.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdcaddremove.sql
new file mode 100644
index 00000000000..e5e59e0e2d4
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdcaddremove.sql
@@ -0,0 +1,206 @@
+-- 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.
+
+-- Define ASNCDC.REMOVETABLE() and ASNCDC.ADDTABLE()
+-- ASNCDC.ADDTABLE() puts a table in CDC mode, making the ASNCapture server collect changes for the table
+-- ASNCDC.REMOVETABLE() makes the ASNCapture server stop collecting changes for that table
+
+--#SET TERMINATOR @
+CREATE OR REPLACE PROCEDURE ASNCDC.REMOVETABLE(
+in tableschema VARCHAR(128),
+in tablename VARCHAR(128)
+)
+LANGUAGE SQL
+P1:
+BEGIN
+
+DECLARE stmtSQL VARCHAR(2048);
+
+DECLARE SQLCODE INT;
+DECLARE SQLSTATE CHAR(5);
+DECLARE RC_SQLCODE INT DEFAULT 0;
+DECLARE RC_SQLSTATE CHAR(5) DEFAULT '00000';
+
+DECLARE CONTINUE HANDLER FOR SQLEXCEPTION, SQLWARNING, NOT FOUND VALUES (SQLCODE, SQLSTATE) INTO RC_SQLCODE, RC_SQLSTATE;
+
+-- delete ASN.IBMSNAP_PRUNCTL entries / source
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMSNAP_PRUNCNTL WHERE SOURCE_OWNER=''' || tableschema || ''' AND SOURCE_TABLE=''' || tablename || '''';
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- delete ASN.IBMSNAP_Register entries / source
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMSNAP_REGISTER WHERE SOURCE_OWNER=''' || tableschema || ''' AND SOURCE_TABLE=''' || tablename || '''';
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- drop CD Table / source
+SET stmtSQL = 'DROP TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ;
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- delete ASN.IBMSNAP_SUBS_COLS entries /target
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMSNAP_SUBS_COLS WHERE TARGET_OWNER=''' || tableschema || ''' AND TARGET_TABLE=''' || tablename || '''';
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- delete ASN.IBMSNAP_SUSBS_MEMBER entries /target
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMSNAP_SUBS_MEMBR WHERE TARGET_OWNER=''' || tableschema || ''' AND TARGET_TABLE=''' || tablename || '''';
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- delete ASN.IBMQREP_COLVERSION
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMQREP_COLVERSION col WHERE EXISTS (SELECT * FROM ASNCDC.IBMQREP_TABVERSION tab WHERE SOURCE_OWNER=''' || tableschema || ''' AND SOURCE_NAME=''' || tablename || '''AND col.TABLEID1 = tab.TABLEID1 AND col.TABLEID2 = tab.TABLEID2';
+ EXECUTE IMMEDIATE stmtSQL;
+
+-- delete ASN.IBMQREP_TABVERSION
+SET stmtSQL = 'DELETE FROM ASNCDC.IBMQREP_TABVERSION WHERE SOURCE_OWNER=''' || tableschema || ''' AND SOURCE_NAME=''' || tablename || '''';
+ EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'ALTER TABLE ' || tableschema || '.' || tablename || ' DATA CAPTURE NONE';
+EXECUTE IMMEDIATE stmtSQL;
+
+END P1@
+--#SET TERMINATOR ;
+
+--#SET TERMINATOR @
+CREATE OR REPLACE PROCEDURE ASNCDC.ADDTABLE(
+in tableschema VARCHAR(128),
+in tablename VARCHAR(128)
+)
+LANGUAGE SQL
+P1:
+BEGIN
+
+DECLARE SQLSTATE CHAR(5);
+
+DECLARE stmtSQL VARCHAR(2048);
+
+SET stmtSQL = 'ALTER TABLE ' || tableschema || '.' || tablename || ' DATA CAPTURE CHANGES';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'CREATE TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' AS ( SELECT ' ||
+ ' CAST('''' AS VARCHAR ( 16 ) FOR BIT DATA) AS IBMSNAP_COMMITSEQ, ' ||
+ ' CAST('''' AS VARCHAR ( 16 ) FOR BIT DATA) AS IBMSNAP_INTENTSEQ, ' ||
+ ' CAST ('''' AS CHAR(1)) ' ||
+ ' AS IBMSNAP_OPERATION, t.* FROM ' || tableschema || '.' || tablename || ' as t ) WITH NO DATA ORGANIZE BY ROW ';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'ALTER TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' ALTER COLUMN IBMSNAP_COMMITSEQ SET NOT NULL';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'ALTER TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' ALTER COLUMN IBMSNAP_INTENTSEQ SET NOT NULL';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'ALTER TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' ALTER COLUMN IBMSNAP_OPERATION SET NOT NULL';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'CREATE UNIQUE INDEX ASNCDC.IXCDC_' ||
+ tableschema || '_' || tablename ||
+ ' ON ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' ( IBMSNAP_COMMITSEQ ASC, IBMSNAP_INTENTSEQ ASC ) PCTFREE 0 MINPCTUSED 0';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'ALTER TABLE ASNCDC.CDC_' ||
+ tableschema || '_' || tablename ||
+ ' VOLATILE CARDINALITY';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'INSERT INTO ASNCDC.IBMSNAP_REGISTER (SOURCE_OWNER, SOURCE_TABLE, ' ||
+ 'SOURCE_VIEW_QUAL, GLOBAL_RECORD, SOURCE_STRUCTURE, SOURCE_CONDENSED, ' ||
+ 'SOURCE_COMPLETE, CD_OWNER, CD_TABLE, PHYS_CHANGE_OWNER, ' ||
+ 'PHYS_CHANGE_TABLE, CD_OLD_SYNCHPOINT, CD_NEW_SYNCHPOINT, ' ||
+ 'DISABLE_REFRESH, CCD_OWNER, CCD_TABLE, CCD_OLD_SYNCHPOINT, ' ||
+ 'SYNCHPOINT, SYNCHTIME, CCD_CONDENSED, CCD_COMPLETE, ARCH_LEVEL, ' ||
+ 'DESCRIPTION, BEFORE_IMG_PREFIX, CONFLICT_LEVEL, ' ||
+ 'CHG_UPD_TO_DEL_INS, CHGONLY, RECAPTURE, OPTION_FLAGS, ' ||
+ 'STOP_ON_ERROR, STATE, STATE_INFO ) VALUES( ' ||
+ '''' || tableschema || ''', ' ||
+ '''' || tablename || ''', ' ||
+ '0, ' ||
+ '''N'', ' ||
+ '1, ' ||
+ '''Y'', ' ||
+ '''Y'', ' ||
+ '''ASNCDC'', ' ||
+ '''CDC_' || tableschema || '_' || tablename || ''', ' ||
+ '''ASNCDC'', ' ||
+ '''CDC_' || tableschema || '_' || tablename || ''', ' ||
+ 'null, ' ||
+ 'null, ' ||
+ '0, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ 'null, ' ||
+ '''0801'', ' ||
+ 'null, ' ||
+ 'null, ' ||
+ '''0'', ' ||
+ '''Y'', ' ||
+ '''N'', ' ||
+ '''Y'', ' ||
+ '''NNNN'', ' ||
+ '''Y'', ' ||
+ '''A'',' ||
+ 'null ) ';
+EXECUTE IMMEDIATE stmtSQL;
+
+SET stmtSQL = 'INSERT INTO ASNCDC.IBMSNAP_PRUNCNTL ( ' ||
+ 'TARGET_SERVER, ' ||
+ 'TARGET_OWNER, ' ||
+ 'TARGET_TABLE, ' ||
+ 'SYNCHTIME, ' ||
+ 'SYNCHPOINT, ' ||
+ 'SOURCE_OWNER, ' ||
+ 'SOURCE_TABLE, ' ||
+ 'SOURCE_VIEW_QUAL, ' ||
+ 'APPLY_QUAL, ' ||
+ 'SET_NAME, ' ||
+ 'CNTL_SERVER , ' ||
+ 'TARGET_STRUCTURE , ' ||
+ 'CNTL_ALIAS , ' ||
+ 'PHYS_CHANGE_OWNER , ' ||
+ 'PHYS_CHANGE_TABLE , ' ||
+ 'MAP_ID ' ||
+ ') VALUES ( ' ||
+ '''KAFKA'', ' ||
+ '''' || tableschema || ''', ' ||
+ '''' || tablename || ''', ' ||
+ 'NULL, ' ||
+ 'NULL, ' ||
+ '''' || tableschema || ''', ' ||
+ '''' || tablename || ''', ' ||
+ '0, ' ||
+ '''KAFKAQUAL'', ' ||
+ '''SET001'', ' ||
+ ' (Select CURRENT_SERVER from sysibm.sysdummy1 ), ' ||
+ '8, ' ||
+ ' (Select CURRENT_SERVER from sysibm.sysdummy1 ), ' ||
+ '''ASNCDC'', ' ||
+ '''CDC_' || tableschema || '_' || tablename || ''', ' ||
+ ' ( SELECT CASE WHEN max(CAST(MAP_ID AS INT)) IS NULL THEN CAST(1 AS VARCHAR(10)) ELSE CAST(CAST(max(MAP_ID) AS INT) + 1 AS VARCHAR(10)) END AS MYINT from ASNCDC.IBMSNAP_PRUNCNTL ) ' ||
+ ' )';
+EXECUTE IMMEDIATE stmtSQL;
+
+END P1@
+--#SET TERMINATOR ;
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdctables.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdctables.sql
new file mode 100644
index 00000000000..45960636b68
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/asncdctables.sql
@@ -0,0 +1,494 @@
+-- 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.
+
+-- 1021 db2 LEVEL Version 10.2.0 --> 11.5.0 1150
+
+CREATE TABLE ASNCDC.IBMQREP_COLVERSION(
+LSN VARCHAR( 16) FOR BIT DATA NOT NULL,
+TABLEID1 SMALLINT NOT NULL,
+TABLEID2 SMALLINT NOT NULL,
+POSITION SMALLINT NOT NULL,
+NAME VARCHAR(128) NOT NULL,
+TYPE SMALLINT NOT NULL,
+LENGTH INTEGER NOT NULL,
+NULLS CHAR( 1) NOT NULL,
+DEFAULT VARCHAR(1536),
+CODEPAGE INTEGER,
+SCALE INTEGER,
+VERSION_TIME TIMESTAMP NOT NULL WITH DEFAULT )
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMQREP_COLVERSIOX
+ON ASNCDC.IBMQREP_COLVERSION(
+LSN ASC,
+TABLEID1 ASC,
+TABLEID2 ASC,
+POSITION ASC);
+
+CREATE INDEX ASNCDC.IX2COLVERSION
+ON ASNCDC.IBMQREP_COLVERSION(
+TABLEID1 ASC,
+TABLEID2 ASC);
+
+CREATE TABLE ASNCDC.IBMQREP_TABVERSION(
+LSN VARCHAR( 16) FOR BIT DATA NOT NULL,
+TABLEID1 SMALLINT NOT NULL,
+TABLEID2 SMALLINT NOT NULL,
+VERSION INTEGER NOT NULL,
+SOURCE_OWNER VARCHAR(128) NOT NULL,
+SOURCE_NAME VARCHAR(128) NOT NULL,
+VERSION_TIME TIMESTAMP NOT NULL WITH DEFAULT )
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMQREP_TABVERSIOX
+ON ASNCDC.IBMQREP_TABVERSION(
+LSN ASC,
+TABLEID1 ASC,
+TABLEID2 ASC,
+VERSION ASC);
+
+CREATE INDEX ASNCDC.IX2TABVERSION
+ON ASNCDC.IBMQREP_TABVERSION(
+TABLEID1 ASC,
+TABLEID2 ASC);
+
+CREATE INDEX ASNCDC.IX3TABVERSION
+ON ASNCDC.IBMQREP_TABVERSION(
+SOURCE_OWNER ASC,
+SOURCE_NAME ASC);
+
+CREATE TABLE ASNCDC.IBMSNAP_APPLEVEL(
+ARCH_LEVEL CHAR( 4) NOT NULL WITH DEFAULT '1021')
+ORGANIZE BY ROW;
+
+INSERT INTO ASNCDC.IBMSNAP_APPLEVEL(ARCH_LEVEL) VALUES (
+'1021');
+
+CREATE TABLE ASNCDC.IBMSNAP_CAPMON(
+MONITOR_TIME TIMESTAMP NOT NULL,
+RESTART_TIME TIMESTAMP NOT NULL,
+CURRENT_MEMORY INT NOT NULL,
+CD_ROWS_INSERTED INT NOT NULL,
+RECAP_ROWS_SKIPPED INT NOT NULL,
+TRIGR_ROWS_SKIPPED INT NOT NULL,
+CHG_ROWS_SKIPPED INT NOT NULL,
+TRANS_PROCESSED INT NOT NULL,
+TRANS_SPILLED INT NOT NULL,
+MAX_TRANS_SIZE INT NOT NULL,
+LOCKING_RETRIES INT NOT NULL,
+JRN_LIB CHAR( 10),
+JRN_NAME CHAR( 10),
+LOGREADLIMIT INT NOT NULL,
+CAPTURE_IDLE INT NOT NULL,
+SYNCHTIME TIMESTAMP NOT NULL,
+CURRENT_LOG_TIME TIMESTAMP NOT NULL WITH DEFAULT ,
+LAST_EOL_TIME TIMESTAMP,
+RESTART_SEQ VARCHAR( 16) FOR BIT DATA NOT NULL WITH DEFAULT ,
+CURRENT_SEQ VARCHAR( 16) FOR BIT DATA NOT NULL WITH DEFAULT ,
+RESTART_MAXCMTSEQ VARCHAR( 16) FOR BIT DATA NOT NULL WITH DEFAULT ,
+LOGREAD_API_TIME INT,
+NUM_LOGREAD_CALLS INT,
+NUM_END_OF_LOGS INT,
+LOGRDR_SLEEPTIME INT,
+NUM_LOGREAD_F_CALLS INT,
+TRANS_QUEUED INT,
+NUM_WARNTXS INT,
+NUM_WARNLOGAPI INT)
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_CAPMONX
+ON ASNCDC.IBMSNAP_CAPMON(
+MONITOR_TIME ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_CAPMON VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_CAPPARMS(
+RETENTION_LIMIT INT,
+LAG_LIMIT INT,
+COMMIT_INTERVAL INT,
+PRUNE_INTERVAL INT,
+TRACE_LIMIT INT,
+MONITOR_LIMIT INT,
+MONITOR_INTERVAL INT,
+MEMORY_LIMIT SMALLINT,
+REMOTE_SRC_SERVER CHAR( 18),
+AUTOPRUNE CHAR( 1),
+TERM CHAR( 1),
+AUTOSTOP CHAR( 1),
+LOGREUSE CHAR( 1),
+LOGSTDOUT CHAR( 1),
+SLEEP_INTERVAL SMALLINT,
+CAPTURE_PATH VARCHAR(1040),
+STARTMODE VARCHAR( 10),
+LOGRDBUFSZ INT NOT NULL WITH DEFAULT 256,
+ARCH_LEVEL CHAR( 4) NOT NULL WITH DEFAULT '1021',
+COMPATIBILITY CHAR( 4) NOT NULL WITH DEFAULT '1021')
+ ORGANIZE BY ROW;
+
+INSERT INTO ASNCDC.IBMSNAP_CAPPARMS(
+RETENTION_LIMIT,
+LAG_LIMIT,
+COMMIT_INTERVAL,
+PRUNE_INTERVAL,
+TRACE_LIMIT,
+MONITOR_LIMIT,
+MONITOR_INTERVAL,
+MEMORY_LIMIT,
+SLEEP_INTERVAL,
+AUTOPRUNE,
+TERM,
+AUTOSTOP,
+LOGREUSE,
+LOGSTDOUT,
+CAPTURE_PATH,
+STARTMODE,
+COMPATIBILITY)
+VALUES (
+10080,
+10080,
+5,
+300,
+10080,
+10080,
+300,
+32,
+5,
+'Y',
+'Y',
+'N',
+'N',
+'N',
+NULL,
+'WARMSI',
+'1021'
+);
+
+CREATE TABLE ASNCDC.IBMSNAP_CAPSCHEMAS (
+ CAP_SCHEMA_NAME VARCHAR(128 OCTETS) NOT NULL
+ )
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_CAPSCHEMASX
+ ON ASNCDC.IBMSNAP_CAPSCHEMAS
+ (CAP_SCHEMA_NAME ASC);
+
+INSERT INTO ASNCDC.IBMSNAP_CAPSCHEMAS(CAP_SCHEMA_NAME) VALUES (
+'ASNCDC');
+
+CREATE TABLE ASNCDC.IBMSNAP_CAPTRACE(
+OPERATION CHAR( 8) NOT NULL,
+TRACE_TIME TIMESTAMP NOT NULL,
+DESCRIPTION VARCHAR(1024) NOT NULL)
+ ORGANIZE BY ROW;
+
+CREATE INDEX ASNCDC.IBMSNAP_CAPTRACEX
+ON ASNCDC.IBMSNAP_CAPTRACE(
+TRACE_TIME ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_CAPTRACE VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_PRUNCNTL(
+TARGET_SERVER CHAR(18) NOT NULL,
+TARGET_OWNER VARCHAR(128) NOT NULL,
+TARGET_TABLE VARCHAR(128) NOT NULL,
+SYNCHTIME TIMESTAMP,
+SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+SOURCE_OWNER VARCHAR(128) NOT NULL,
+SOURCE_TABLE VARCHAR(128) NOT NULL,
+SOURCE_VIEW_QUAL SMALLINT NOT NULL,
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+CNTL_SERVER CHAR( 18) NOT NULL,
+TARGET_STRUCTURE SMALLINT NOT NULL,
+CNTL_ALIAS CHAR( 8),
+PHYS_CHANGE_OWNER VARCHAR(128),
+PHYS_CHANGE_TABLE VARCHAR(128),
+MAP_ID VARCHAR(10) NOT NULL)
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_PRUNCNTLX
+ON ASNCDC.IBMSNAP_PRUNCNTL(
+SOURCE_OWNER ASC,
+SOURCE_TABLE ASC,
+SOURCE_VIEW_QUAL ASC,
+APPLY_QUAL ASC,
+SET_NAME ASC,
+TARGET_SERVER ASC,
+TARGET_TABLE ASC,
+TARGET_OWNER ASC);
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_PRUNCNTLX1
+ON ASNCDC.IBMSNAP_PRUNCNTL(
+MAP_ID ASC);
+
+CREATE INDEX ASNCDC.IBMSNAP_PRUNCNTLX2
+ON ASNCDC.IBMSNAP_PRUNCNTL(
+PHYS_CHANGE_OWNER ASC,
+PHYS_CHANGE_TABLE ASC);
+
+CREATE INDEX ASNCDC.IBMSNAP_PRUNCNTLX3
+ON ASNCDC.IBMSNAP_PRUNCNTL(
+APPLY_QUAL ASC,
+SET_NAME ASC,
+TARGET_SERVER ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_PRUNCNTL VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_PRUNE_LOCK(
+DUMMY CHAR( 1))
+ ORGANIZE BY ROW;
+
+CREATE TABLE ASNCDC.IBMSNAP_PRUNE_SET(
+TARGET_SERVER CHAR( 18) NOT NULL,
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+SYNCHTIME TIMESTAMP,
+SYNCHPOINT VARCHAR( 16) FOR BIT DATA NOT NULL)
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_PRUNE_SETX
+ON ASNCDC.IBMSNAP_PRUNE_SET(
+TARGET_SERVER ASC,
+APPLY_QUAL ASC,
+SET_NAME ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_PRUNE_SET VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_REGISTER(
+SOURCE_OWNER VARCHAR(128) NOT NULL,
+SOURCE_TABLE VARCHAR(128) NOT NULL,
+SOURCE_VIEW_QUAL SMALLINT NOT NULL,
+GLOBAL_RECORD CHAR( 1) NOT NULL,
+SOURCE_STRUCTURE SMALLINT NOT NULL,
+SOURCE_CONDENSED CHAR( 1) NOT NULL,
+SOURCE_COMPLETE CHAR( 1) NOT NULL,
+CD_OWNER VARCHAR(128),
+CD_TABLE VARCHAR(128),
+PHYS_CHANGE_OWNER VARCHAR(128),
+PHYS_CHANGE_TABLE VARCHAR(128),
+CD_OLD_SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+CD_NEW_SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+DISABLE_REFRESH SMALLINT NOT NULL,
+CCD_OWNER VARCHAR(128),
+CCD_TABLE VARCHAR(128),
+CCD_OLD_SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+SYNCHTIME TIMESTAMP,
+CCD_CONDENSED CHAR( 1),
+CCD_COMPLETE CHAR( 1),
+ARCH_LEVEL CHAR( 4) NOT NULL,
+DESCRIPTION CHAR(254),
+BEFORE_IMG_PREFIX VARCHAR( 4),
+CONFLICT_LEVEL CHAR( 1),
+CHG_UPD_TO_DEL_INS CHAR( 1),
+CHGONLY CHAR( 1),
+RECAPTURE CHAR( 1),
+OPTION_FLAGS CHAR( 4) NOT NULL,
+STOP_ON_ERROR CHAR( 1) WITH DEFAULT 'Y',
+STATE CHAR( 1) WITH DEFAULT 'I',
+STATE_INFO CHAR( 8))
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_REGISTERX
+ON ASNCDC.IBMSNAP_REGISTER(
+SOURCE_OWNER ASC,
+SOURCE_TABLE ASC,
+SOURCE_VIEW_QUAL ASC);
+
+CREATE INDEX ASNCDC.IBMSNAP_REGISTERX1
+ON ASNCDC.IBMSNAP_REGISTER(
+PHYS_CHANGE_OWNER ASC,
+PHYS_CHANGE_TABLE ASC);
+
+CREATE INDEX ASNCDC.IBMSNAP_REGISTERX2
+ON ASNCDC.IBMSNAP_REGISTER(
+GLOBAL_RECORD ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_REGISTER VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_RESTART(
+MAX_COMMITSEQ VARCHAR( 16) FOR BIT DATA NOT NULL,
+MAX_COMMIT_TIME TIMESTAMP NOT NULL,
+MIN_INFLIGHTSEQ VARCHAR( 16) FOR BIT DATA NOT NULL,
+CURR_COMMIT_TIME TIMESTAMP NOT NULL,
+CAPTURE_FIRST_SEQ VARCHAR( 16) FOR BIT DATA NOT NULL)
+ ORGANIZE BY ROW;
+
+CREATE TABLE ASNCDC.IBMSNAP_SIGNAL(
+SIGNAL_TIME TIMESTAMP NOT NULL WITH DEFAULT ,
+SIGNAL_TYPE VARCHAR( 30) NOT NULL,
+SIGNAL_SUBTYPE VARCHAR( 30),
+SIGNAL_INPUT_IN VARCHAR(500),
+SIGNAL_STATE CHAR( 1) NOT NULL,
+SIGNAL_LSN VARCHAR( 16) FOR BIT DATA)
+DATA CAPTURE CHANGES
+ ORGANIZE BY ROW;
+
+CREATE INDEX ASNCDC.IBMSNAP_SIGNALX
+ON ASNCDC.IBMSNAP_SIGNAL(
+SIGNAL_TIME ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_SIGNAL VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_SUBS_COLS(
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+WHOS_ON_FIRST CHAR( 1) NOT NULL,
+TARGET_OWNER VARCHAR(128) NOT NULL,
+TARGET_TABLE VARCHAR(128) NOT NULL,
+COL_TYPE CHAR( 1) NOT NULL,
+TARGET_NAME VARCHAR(128) NOT NULL,
+IS_KEY CHAR( 1) NOT NULL,
+COLNO SMALLINT NOT NULL,
+EXPRESSION VARCHAR(1024) NOT NULL)
+ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_SUBS_COLSX
+ON ASNCDC.IBMSNAP_SUBS_COLS(
+APPLY_QUAL ASC,
+SET_NAME ASC,
+WHOS_ON_FIRST ASC,
+TARGET_OWNER ASC,
+TARGET_TABLE ASC,
+TARGET_NAME ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_SUBS_COLS VOLATILE CARDINALITY;
+
+--CREATE UNIQUE INDEX ASNCDC.IBMSNAP_SUBS_EVENTX
+--ON ASNCDC.IBMSNAP_SUBS_EVENT(
+--EVENT_NAME ASC,
+--EVENT_TIME ASC);
+
+
+--ALTER TABLE ASNCDC.IBMSNAP_SUBS_EVENT VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_SUBS_MEMBR(
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+WHOS_ON_FIRST CHAR( 1) NOT NULL,
+SOURCE_OWNER VARCHAR(128) NOT NULL,
+SOURCE_TABLE VARCHAR(128) NOT NULL,
+SOURCE_VIEW_QUAL SMALLINT NOT NULL,
+TARGET_OWNER VARCHAR(128) NOT NULL,
+TARGET_TABLE VARCHAR(128) NOT NULL,
+TARGET_CONDENSED CHAR( 1) NOT NULL,
+TARGET_COMPLETE CHAR( 1) NOT NULL,
+TARGET_STRUCTURE SMALLINT NOT NULL,
+PREDICATES VARCHAR(1024),
+MEMBER_STATE CHAR( 1),
+TARGET_KEY_CHG CHAR( 1) NOT NULL,
+UOW_CD_PREDICATES VARCHAR(1024),
+JOIN_UOW_CD CHAR( 1),
+LOADX_TYPE SMALLINT,
+LOADX_SRC_N_OWNER VARCHAR( 128),
+LOADX_SRC_N_TABLE VARCHAR(128))
+ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_SUBS_MEMBRX
+ON ASNCDC.IBMSNAP_SUBS_MEMBR(
+APPLY_QUAL ASC,
+SET_NAME ASC,
+WHOS_ON_FIRST ASC,
+SOURCE_OWNER ASC,
+SOURCE_TABLE ASC,
+SOURCE_VIEW_QUAL ASC,
+TARGET_OWNER ASC,
+TARGET_TABLE ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_SUBS_MEMBR VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_SUBS_SET(
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+SET_TYPE CHAR( 1) NOT NULL,
+WHOS_ON_FIRST CHAR( 1) NOT NULL,
+ACTIVATE SMALLINT NOT NULL,
+SOURCE_SERVER CHAR( 18) NOT NULL,
+SOURCE_ALIAS CHAR( 8),
+TARGET_SERVER CHAR( 18) NOT NULL,
+TARGET_ALIAS CHAR( 8),
+STATUS SMALLINT NOT NULL,
+LASTRUN TIMESTAMP NOT NULL,
+REFRESH_TYPE CHAR( 1) NOT NULL,
+SLEEP_MINUTES INT,
+EVENT_NAME CHAR( 18),
+LASTSUCCESS TIMESTAMP,
+SYNCHPOINT VARCHAR( 16) FOR BIT DATA,
+SYNCHTIME TIMESTAMP,
+CAPTURE_SCHEMA VARCHAR(128) NOT NULL,
+TGT_CAPTURE_SCHEMA VARCHAR(128),
+FEDERATED_SRC_SRVR VARCHAR( 18),
+FEDERATED_TGT_SRVR VARCHAR( 18),
+JRN_LIB CHAR( 10),
+JRN_NAME CHAR( 10),
+OPTION_FLAGS CHAR( 4) NOT NULL,
+COMMIT_COUNT SMALLINT,
+MAX_SYNCH_MINUTES SMALLINT,
+AUX_STMTS SMALLINT NOT NULL,
+ARCH_LEVEL CHAR( 4) NOT NULL)
+ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_SUBS_SETX
+ON ASNCDC.IBMSNAP_SUBS_SET(
+APPLY_QUAL ASC,
+SET_NAME ASC,
+WHOS_ON_FIRST ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_SUBS_SET VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_SUBS_STMTS(
+APPLY_QUAL CHAR( 18) NOT NULL,
+SET_NAME CHAR( 18) NOT NULL,
+WHOS_ON_FIRST CHAR( 1) NOT NULL,
+BEFORE_OR_AFTER CHAR( 1) NOT NULL,
+STMT_NUMBER SMALLINT NOT NULL,
+EI_OR_CALL CHAR( 1) NOT NULL,
+SQL_STMT VARCHAR(1024),
+ACCEPT_SQLSTATES VARCHAR( 50))
+ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_SUBS_STMTSX
+ON ASNCDC.IBMSNAP_SUBS_STMTS(
+APPLY_QUAL ASC,
+SET_NAME ASC,
+WHOS_ON_FIRST ASC,
+BEFORE_OR_AFTER ASC,
+STMT_NUMBER ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_SUBS_STMTS VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_UOW(
+IBMSNAP_UOWID CHAR( 10) FOR BIT DATA NOT NULL,
+IBMSNAP_COMMITSEQ VARCHAR( 16) FOR BIT DATA NOT NULL,
+IBMSNAP_LOGMARKER TIMESTAMP NOT NULL,
+IBMSNAP_AUTHTKN VARCHAR(30) NOT NULL,
+IBMSNAP_AUTHID VARCHAR(128) NOT NULL,
+IBMSNAP_REJ_CODE CHAR( 1) NOT NULL WITH DEFAULT ,
+IBMSNAP_APPLY_QUAL CHAR( 18) NOT NULL WITH DEFAULT )
+ ORGANIZE BY ROW;
+
+CREATE UNIQUE INDEX ASNCDC.IBMSNAP_UOWX
+ON ASNCDC.IBMSNAP_UOW(
+IBMSNAP_COMMITSEQ ASC,
+IBMSNAP_LOGMARKER ASC);
+
+ALTER TABLE ASNCDC.IBMSNAP_UOW VOLATILE CARDINALITY;
+
+CREATE TABLE ASNCDC.IBMSNAP_CAPENQ (
+ LOCK_NAME CHAR(9 OCTETS)
+ )
+ ORGANIZE BY ROW
+ DATA CAPTURE NONE
+ COMPRESS NO;
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/cdcsetup.sh b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/cdcsetup.sh
new file mode 100644
index 00000000000..c943790fd05
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/cdcsetup.sh
@@ -0,0 +1,34 @@
+#/bin/bash
+################################################################################
+# 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.
+################################################################################
+
+if [ ! -f /asncdctools/src/asncdc.nlk ]; then
+rc=1
+echo "waiting for db2inst1 exists ."
+while [ "$rc" -ne 0 ]
+do
+ sleep 5
+ id db2inst1
+ rc=$?
+ echo '.'
+done
+
+su -c "/asncdctools/src/dbsetup.sh $DBNAME" - db2inst1
+fi
+touch /asncdctools/src/asncdc.nlk
+
+echo "The asncdc program enable finished"
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/column_type_test.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/column_type_test.sql
new file mode 100644
index 00000000000..b21e33ed583
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/column_type_test.sql
@@ -0,0 +1,49 @@
+-- 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.
+
+-- ----------------------------------------------------------------------------------------------------------------
+-- DATABASE: column_type_test
+-- ----------------------------------------------------------------------------------------------------------------
+
+CREATE TABLE DB2INST1.FULL_TYPES (
+ ID INTEGER NOT NULL,
+ SMALL_C SMALLINT,
+ INT_C INTEGER,
+ BIG_C BIGINT,
+ REAL_C REAL,
+ DOUBLE_C DOUBLE,
+ NUMERIC_C NUMERIC(10, 5),
+ DECIMAL_C DECIMAL(10, 1),
+ VARCHAR_C VARCHAR(200),
+ CHAR_C CHAR,
+ CHARACTER_C CHAR(3),
+ TIMESTAMP_C TIMESTAMP,
+ DATE_C DATE,
+ TIME_C TIME,
+ DEFAULT_NUMERIC_C NUMERIC,
+ TIMESTAMP_PRECISION_C TIMESTAMP(9),
+ PRIMARY KEY (ID)
+);
+
+INSERT INTO DB2INST1.FULL_TYPES VALUES (
+ 1, 32767, 65535, 2147483647, 5.5, 6.6, 123.12345, 404.4443,
+ 'Hello World', 'a', 'abc', '2020-07-17 18:00:22.123', '2020-07-17', '18:00:22', 500,
+ '2020-07-17 18:00:22.123456789');
+
+VALUES ASNCDC.ASNCDCSERVICES('status','asncdc');
+
+CALL ASNCDC.ADDTABLE('DB2INST1', 'FULL_TYPES');
+
+VALUES ASNCDC.ASNCDCSERVICES('reinit','asncdc');
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/customers.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/customers.sql
new file mode 100644
index 00000000000..3499788b1d5
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/customers.sql
@@ -0,0 +1,50 @@
+-- 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.
+
+CREATE TABLE DB2INST1.CUSTOMERS (
+ ID INTEGER NOT NULL PRIMARY KEY,
+ NAME VARCHAR(255) NOT NULL DEFAULT 'flink',
+ ADDRESS VARCHAR(1024),
+ PHONE_NUMBER VARCHAR(512)
+);
+
+INSERT INTO DB2INST1.CUSTOMERS
+VALUES (101,'user_1','Shanghai','123567891234'),
+ (102,'user_2','Shanghai','123567891234'),
+ (103,'user_3','Shanghai','123567891234'),
+ (109,'user_4','Shanghai','123567891234'),
+ (110,'user_5','Shanghai','123567891234'),
+ (111,'user_6','Shanghai','123567891234'),
+ (118,'user_7','Shanghai','123567891234'),
+ (121,'user_8','Shanghai','123567891234'),
+ (123,'user_9','Shanghai','123567891234'),
+ (1009,'user_10','Shanghai','123567891234'),
+ (1010,'user_11','Shanghai','123567891234'),
+ (1011,'user_12','Shanghai','123567891234'),
+ (1012,'user_13','Shanghai','123567891234'),
+ (1013,'user_14','Shanghai','123567891234'),
+ (1014,'user_15','Shanghai','123567891234'),
+ (1015,'user_16','Shanghai','123567891234'),
+ (1016,'user_17','Shanghai','123567891234'),
+ (1017,'user_18','Shanghai','123567891234'),
+ (1018,'user_19','Shanghai','123567891234'),
+ (1019,'user_20','Shanghai','123567891234'),
+ (2000,'user_21','Shanghai','123567891234');
+
+VALUES ASNCDC.ASNCDCSERVICES('status','asncdc');
+
+CALL ASNCDC.ADDTABLE('DB2INST1', 'CUSTOMERS');
+
+VALUES ASNCDC.ASNCDCSERVICES('reinit','asncdc');
\ No newline at end of file
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/dbsetup.sh b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/dbsetup.sh
new file mode 100644
index 00000000000..44ad8e5d9f3
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/dbsetup.sh
@@ -0,0 +1,67 @@
+#/bin/bash
+################################################################################
+# 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.
+################################################################################
+
+echo "Compile ASN tool ..."
+cd /asncdctools/src
+/opt/ibm/db2/V11.5/samples/c/bldrtn asncdc
+
+DBNAME=$1
+DB2DIR=/opt/ibm/db2/V11.5
+rc=1
+echo "Waiting for DB2 start ( $DBNAME ) ."
+while [ "$rc" -ne 0 ]
+do
+ sleep 5
+ db2 connect to $DBNAME
+ rc=$?
+ echo '.'
+done
+
+# enable metacatalog read via JDBC
+cd $HOME/sqllib/bnd
+db2 bind db2schema.bnd blocking all grant public sqlerror continue
+
+# do a backup and restart the db
+db2 backup db $DBNAME to /dev/null
+db2 restart db $DBNAME
+
+db2 connect to $DBNAME
+
+cp /asncdctools/src/asncdc /database/config/db2inst1/sqllib/function
+chmod 777 /database/config/db2inst1/sqllib/function
+
+# add UDF / start stop asncap
+db2 -tvmf /asncdctools/src/asncdc_UDF.sql
+
+# create asntables
+db2 -tvmf /asncdctools/src/asncdctables.sql
+
+# add UDF / add remove asntables
+
+db2 -tvmf /asncdctools/src/asncdcaddremove.sql
+
+
+
+
+# startup-agent
+db2 -tvmf /asncdctools/src/startup-agent.sql
+
+
+
+
+echo "db2 setup done"
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/inventory.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/inventory.sql
new file mode 100644
index 00000000000..a665d96fc57
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/inventory.sql
@@ -0,0 +1,40 @@
+-- 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.
+
+-- Create and populate our test products table using a single insert with many rows
+CREATE TABLE DB2INST1.PRODUCTS (
+ ID INTEGER NOT NULL GENERATED BY DEFAULT AS IDENTITY
+ (START WITH 101, INCREMENT BY 1) PRIMARY KEY,
+ NAME VARCHAR(255) NOT NULL,
+ DESCRIPTION VARCHAR(512),
+ WEIGHT FLOAT
+);
+
+INSERT INTO DB2INST1.PRODUCTS(NAME,DESCRIPTION,WEIGHT)
+VALUES ('scooter','Small 2-wheel scooter',3.14),
+ ('car battery','12V car battery',8.1),
+ ('12-pack drill bits','12-pack of drill bits with sizes ranging from #40 to #3',0.8),
+ ('hammer','12oz carpenter''s hammer',0.75),
+ ('hammer','14oz carpenter''s hammer',0.875),
+ ('hammer','16oz carpenter''s hammer',1.0),
+ ('rocks','box of assorted rocks',5.3),
+ ('jacket','water resistent black wind breaker',0.1),
+ ('spare tire','24 inch spare tire',22.2);
+
+VALUES ASNCDC.ASNCDCSERVICES('status','asncdc');
+
+CALL ASNCDC.ADDTABLE('DB2INST1', 'PRODUCTS');
+
+VALUES ASNCDC.ASNCDCSERVICES('reinit','asncdc');
\ No newline at end of file
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/startup-agent.sql b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/startup-agent.sql
new file mode 100644
index 00000000000..4288b7bb339
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/db2_server/startup-agent.sql
@@ -0,0 +1,16 @@
+-- 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.
+
+VALUES ASNCDC.ASNCDCSERVICES('start','asncdc');
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/log4j2-test.properties b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/log4j2-test.properties
new file mode 100644
index 00000000000..32df1c0251c
--- /dev/null
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-db2/src/test/resources/log4j2-test.properties
@@ -0,0 +1,25 @@
+################################################################################
+# 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.
+################################################################################
+
+# Set root logger level to ERROR to not flood build logs
+# set manually to INFO for debugging purposes
+rootLogger.level=ERROR
+rootLogger.appenderRef.test.ref = TestLogger
+
+appender.testlogger.name = TestLogger
+appender.testlogger.type = CONSOLE
+appender.testlogger.target = SYSTEM_ERR
+appender.testlogger.layout.type = PatternLayout
+appender.testlogger.layout.pattern = %-4r [%t] %-5p %c - %m%n
diff --git a/flink-cdc-connect/flink-cdc-pipeline-connectors/pom.xml b/flink-cdc-connect/flink-cdc-pipeline-connectors/pom.xml
index 490a0d33a6c..006a362acea 100644
--- a/flink-cdc-connect/flink-cdc-pipeline-connectors/pom.xml
+++ b/flink-cdc-connect/flink-cdc-pipeline-connectors/pom.xml
@@ -43,6 +43,7 @@ limitations under the License.
flink-cdc-pipeline-connector-fluss
flink-cdc-pipeline-connector-sqlserver
flink-cdc-pipeline-connector-hudi
+ flink-cdc-pipeline-connector-db2