改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

FlinkCDC实时数据同步实战指南

ckckck2025年1月10日6 浏览

Flink CDC 整库同步 DataStream API

场景一:MySQL 整库同步到 MySQL

思路:

  1. 使用 mysqlcatalog 获取各个表的信息(列名、列类型等)。
  2. 创建相应的 Sink Table。
  3. Flink CDC 的 DataStream 提供了整库获取数据的能力,因此我们采用 DataStream 的方式获取数据。
  4. 在自定义反序列化中形成 <tableName, Row> 的输出,得到 DataStream<Tuple2<String, Row>>
  5. 根据 tableName 将流拆分(过滤),相当于每个 tableName 对应一个自己的 DataStream
  6. 将每个流转换为 SourceTable,然后执行 INSERT INTO SinkTable SELECT * FROM SourceTable

pom.xml 文件

xml 复制代码
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>
    <groupId>com.cityos</groupId>
    <artifactId>flink_1_15</artifactId>
    <version>1.0-SNAPSHOT</version>
    <properties>
        <java.version>1.8</java.version>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
        <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
        <spring-boot.version>2.3.7.RELEASE</spring-boot.version>
        <flink.version>1.15.2</flink.version>
        <scala.binary.version>2.12</scala.binary.version>
    </properties>
    <repositories>
        <repository>
            <id>scala-tools.org</id>
            <name>Scala-Tools Maven2 Repository</name>
            <url>http://scala-tools.org/repo-releases</url>
        </repository>
        <repository>
            <id>spring</id>
            <url>https://maven.aliyun.com/repository/spring</url>
        </repository>
        <repository>
            <id>cloudera</id>
            <url>https://repository.cloudera.com/artifactory/cloudera-repos/</url>
        </repository>
    </repositories>
    <dependencies>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-scala_${scala.binary.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-scala_${scala.binary.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-planner_${scala.binary.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-api-scala-bridge_${scala.binary.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-common</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-kafka</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-connector-jdbc</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>com.ververica</groupId>
            <artifactId>flink-connector-mysql-cdc</artifactId>
            <version>2.3.0</version>
        </dependency>
        <dependency>
            <groupId>mysql</groupId>
            <artifactId>mysql-connector-java</artifactId>
            <version>8.0.29</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-json</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-csv</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.slf4j</groupId>
            <artifactId>slf4j-log4j12</artifactId>
            <version>1.7.21</version>
            <scope>compile</scope>
        </dependency>
        <dependency>
            <groupId>log4j</groupId>
            <artifactId>log4j</artifactId>
            <version>1.2.17</version>
        </dependency>
    </dependencies>
    <dependencyManagement>
        <dependencies>
            <dependency>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-dependencies</artifactId>
                <version>${spring-boot.version}</version>
                <type>pom</type>
                <scope>import</scope>
            </dependency>
        </dependencies>
    </dependencyManagement>
    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.8.1</version>
                <configuration>
                    <source>1.8</source>
                    <target>1.8</target>
                    <encoding>UTF-8</encoding>
                </configuration>
            </plugin>
            <plugin>
                <groupId>org.springframework.boot</groupId>
                <artifactId>spring-boot-maven-plugin</artifactId>
                <version>2.3.7.RELEASE</version>
                <configuration>
                    <mainClass>com.cityos.Flink1142Application</mainClass>
                </configuration>
                <executions>
                    <execution>
                        <id>repackage</id>
                        <goals>
                            <goal>repackage</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

同步代码

java 复制代码
package com.cityos;

import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import org.apache.commons.lang3.StringUtils;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.calcite.shaded.com.google.common.collect.Maps;
import org.apache.flink.connector.jdbc.catalog.MySqlCatalog;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Schema;
import org.apache.flink.table.api.StatementSet;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.catalog.DefaultCatalogTable;
import org.apache.flink.table.catalog.ObjectPath;
import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.table.types.logical.RowType;
import org.apache.flink.types.Row;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.ArrayList;
import java.util.List;
import java.util.Map;

public class FlinkCdcMultiSyncJdbc {

    private static final Logger log = LoggerFactory.getLogger(FlinkCdcMultiSyncJdbc.class);

    public static void main(String[] args) throws Exception {
        // 解析参数
        ParameterTool parameterTool = ParameterTool.fromArgs(args);
        String userName = parameterTool.get("username");
        String passWord = parameterTool.get("password");
        String host = parameterTool.get("host");
        String db = parameterTool.get("db");
        int port = Integer.valueOf(Optional.ofNullable(parameterTool.get("port")).orElse("3306"));

        // 读取传入参数文件,获取 FlinkSQL DDL WITH 参数 SQL
        String connectorWithBodyPath = parameterTool.get("mysql_with_body_path");
        List<String> connectorWithBodyList = Files.readAllLines(Paths.get(connectorWithBodyPath));
        String connectorWithBody = StringUtils.join(connectorWithBodyList, "
");

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.enableCheckpointing(3000);
        StreamTableEnvironment tEnv = StreamTableEnvironment.create(env);

        // 注册同步的库对应的 Catalog
        MySqlCatalog mysqlCatalog = new MySqlCatalog("mysql-catalog", db, userName, passWord, String.format("jdbc:mysql://%s:%d", host, port));
        List<String> tables = new ArrayList<>();

        // 如果整库同步,则从 Catalog 里取所有表,否则从指定表中取表名
        if (".*".equals(tableList)) {
            tables = mysqlCatalog.listTables(db);
        } else {
            String[] tableArray = tableList.split(",");
            for (String table : tableArray) {
                tables.add(table.split("\.")[1]);
            }
        }

        // 创建表名和对应 RowTypeInfo 映射的 Map
        Map<String, RowTypeInfo> tableTypeInformationMap = Maps.newConcurrentMap();
        Map<String, DataType[]> tableDataTypesMap = Maps.newConcurrentMap();
        Map<String, RowType> tableRowTypeMap = Maps.newConcurrentMap();
        for (String table : tables) {
            // 获取 MySQL Catalog 中注册的表
            ObjectPath objectPath = new ObjectPath(db, table);
            DefaultCatalogTable catalogBaseTable = (DefaultCatalogTable) mysqlCatalog.getTable(objectPath);
            // 获取表的 Schema
            Schema schema = catalogBaseTable.getUnresolvedSchema();
            // 获取表中字段名列表
            String[] fieldNames = new String[schema.getColumns().size()];
            // 获取 DataType
            DataType[] fieldDataTypes = new DataType[schema.getColumns().size()];
            LogicalType[] logicalTypes = new LogicalType[schema.getColumns().size()];
            // 获取表字段类型
            TypeInformation<?>[] fieldTypes = new TypeInformation[schema.getColumns().size()];
            // 获取表的主键
            List<String> primaryKeys = schema.getPrimaryKey().get().getColumnNames();

            for (int i = 0; i < schema.getColumns().size(); i++) {
                Schema.UnresolvedPhysicalColumn column = (Schema.UnresolvedPhysicalColumn) schema.getColumns().get(i);
                fieldNames[i] = column.getName();
                fieldDataTypes[i] = (DataType) column.getDataType();
                fieldTypes[i] = InternalTypeInfo.of(((DataType) column.getDataType()).getLogicalType());
                logicalTypes[i] = ((DataType) column.getDataType()).getLogicalType();
            }
            RowType rowType = RowType.of(logicalTypes, fieldNames);
            tableRowTypeMap.put(table, rowType);

            // 组装 Sink 表 DDL SQL
            StringBuilder stmt = new StringBuilder();
            String tableName = table;
            String jdbcSinkTableName = String.format("jdbc_sink_%s", tableName);
            stmt.append("create table ").append(jdbcSinkTableName).append("(
");

            for (int i = 0; i < fieldNames.length; i++) {
                String column = fieldNames[i];
                String fieldDataType = fieldDataTypes[i].toString();
                stmt.append("	").append(column).append(" ").append(fieldDataType).append(",
");
            }
            stmt.append(String.format("PRIMARY KEY (%s) NOT ENFORCED
)", StringUtils.join(primaryKeys, ",")));
            String formatJdbcSinkWithBody = connectorWithBody
                    .replace("${tableName}", jdbcSinkTableName);
            String createSinkTableDdl = stmt.toString() + formatJdbcSinkWithBody;
            // 创建 Sink 表
            log.info("createSinkTableDdl: {}", createSinkTableDdl);
            tEnv.executeSql(createSinkTableDdl);
            tableDataTypesMap.put(tableName, fieldDataTypes);
            tableTypeInformationMap.put(tableName, new RowTypeInfo(fieldTypes, fieldNames));
        }

        // 监控 MySQL Binlog
        Properties dbzProperties = new Properties();
        dbzProperties.setProperty("decimal.handling.mode", "string");
        dbzProperties.setProperty("converters", "data");
        dbzProperties.setProperty("data.type", "com.flink.flinksql.flinkcdc.CustomerDataConverter");
        dbzProperties.setProperty("data.format.date", "yyyy-MM-dd");
        dbzProperties.setProperty("data.format.time", "HH:mm:ss");
        dbzProperties.setProperty("data.format.datetime", "yyyy-MM-dd HH:mm:ss");
        dbzProperties.setProperty("data.format.timestamp", "yyyy-MM-dd HH:mm:ss");
        dbzProperties.setProperty("data.format.timestamp.zone", "UTC+8");

        MySqlSource<Tuple2<String, Row>> mySqlSource = MySqlSource.<Tuple2<String, Row>>builder()
                .debeziumProperties(dbzProperties)
                .hostname(host)
                .port(port)
                .databaseList(db)
                .tableList(String.format("%s.*", db))
                .username(userName)
                .password(passWord)
                .deserializer(new CustomDebeziumDeserializer(tableRowTypeMap))
                .startupOptions(StartupOptions.initial())
                .build();

        SingleOutputStreamOperator<Tuple2<String, Row>> dataStreamSource = env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "mysql cdc").disableChaining();
        StatementSet statementSet = tEnv.createStatementSet();

        // dataStream 转 Table,创建临时视图,插入 Sink 表
        for (Map.Entry<String, RowTypeInfo> entry : tableTypeInformationMap.entrySet()) {
            String tableName = entry.getKey();
            RowTypeInfo rowTypeInfo = entry.getValue();
            SingleOutputStreamOperator<Row> mapStream = dataStreamSource.filter(data -> data.f0.equals(tableName)).map(data -> data.f1, rowTypeInfo);
            Table table = tEnv.fromChangelogStream(mapStream);
            String temporaryViewName = String.format("t_%s", tableName);
            tEnv.createTemporaryView(temporaryViewName, table);
            String sinkTableName = String.format("jdbc_sink_%s", tableName);
            String insertSql = String.format("insert into %s select * from %s", sinkTableName, temporaryViewName);
            log.info("add insertSql for {},sql: {}", tableName, insertSql);
            statementSet.addInsertSql(insertSql);
        }
        statementSet.execute();
    }
}

反序列化

java 复制代码
package com.cityos;

import com.ververica.cdc.debezium.DebeziumDeserializationSchema;
import com.ververica.cdc.debezium.table.DeserializationRuntimeConverter;
import com.ververica.cdc.debezium.utils.TemporalConversions;
import io.debezium.data.Envelope;
import io.debezium.data.SpecialValueDecimal;
import io.debezium.data.VariableScaleDecimal;
import io.debezium.time.*;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.calcite.shaded.com.google.common.collect.Maps;
import org.apache.flink.table.data.DecimalData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.StringData;
import org.apache.flink.table.data.TimestampData;
import org.apache.flink.table.types.logical.DecimalType;
import org.apache.flink.table.types.logical.LogicalType;
import org.apache.flink.table.types.logical.RowType;
import org.apache.flink.types.Row;
import org.apache.flink.types.RowKind;
import org.apache.flink.util.Collector;
import org.apache.kafka.connect.data.Decimal;
import org.apache.kafka.connect.data.Field;
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.data.Struct;
import org.apache.kafka.connect.source.SourceRecord;

import java.math.BigDecimal;
import java.nio.ByteBuffer;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.Map;

public class CustomDebeziumDeserializer implements DebeziumDeserializationSchema<Tuple2<String, Row>> {

    private final Map<String, RowType> tableRowTypeMap;
    private Map<String, DeserializationRuntimeConverter> physicalConverterMap = Maps.newConcurrentMap();

    CustomDebeziumDeserializer(Map<String, RowType> tableRowTypeMap) {
        this.tableRowTypeMap = tableRowTypeMap;
        for (String tablename : this.tableRowTypeMap.keySet()) {
            RowType rowType = this.tableRowTypeMap.get(tablename);
            DeserializationRuntimeConverter physicalConverter = createNotNullConverter(rowType);
            this.physicalConverterMap.put(tablename, physicalConverter);
        }
    }

    @Override
    public void deserialize(SourceRecord record, Collector<Tuple2<String, Row>> out) throws Exception {
        Envelope.Operation op = Envelope.operationFor(record);
        Struct value = (Struct) record.value();
        Schema valueSchema = record.valueSchema();
        Struct source = value.getStruct("source");
        String tablename = source.get("table").toString();
        DeserializationRuntimeConverter physicalConverter = physicalConverterMap.get(tablename);
        if (op == Envelope.Operation.CREATE || op == Envelope.Operation.READ) {
            Row insert = extractAfterRow(value, valueSchema, physicalConverter);
            insert.setKind(RowKind.INSERT);
            out.collect(Tuple2.of(tablename, insert));
        } else if (op == Envelope.Operation.DELETE) {
            Row delete = extractBeforeRow(value, valueSchema, physicalConverter);
            delete.setKind(RowKind.DELETE);
            out.collect(Tuple2.of(tablename, delete));
        } else {
            Row before = extractBeforeRow(value, valueSchema, physicalConverter);
            before.setKind(RowKind.UPDATE_BEFORE);
            out.collect(Tuple2.of(tablename, before));

            Row after = extractAfterRow(value, valueSchema, physicalConverter);
            after.setKind(RowKind.UPDATE_AFTER);
            out.collect(Tuple2.of(tablename, after));
        }
    }

    private Row extractAfterRow(Struct value, Schema valueSchema, DeserializationRuntimeConverter physicalConverter) throws Exception {
        Schema afterSchema = valueSchema.field(Envelope.FieldName.AFTER).schema();
        Struct after = value.getStruct(Envelope.FieldName.AFTER);
        return (Row) physicalConverter.convert(after, afterSchema);
    }

    private Row extractBeforeRow(Struct value, Schema valueSchema, DeserializationRuntimeConverter physicalConverter) throws Exception {
        Schema beforeSchema = valueSchema.field(Envelope.FieldName.BEFORE).schema();
        Struct before = value.getStruct(Envelope.FieldName.BEFORE);
        return (Row) physicalConverter.convert(before, beforeSchema);
    }

    @Override
    public TypeInformation<Tuple2<String, Row>> getProducedType() {
        return TypeInformation.of(new TypeHint<Tuple2<String, Row>>() {
        });
    }

    public static DeserializationRuntimeConverter createNotNullConverter(LogicalType type) {
        switch (type.getTypeRoot()) {
            case NULL:
                return (dbzObj, schema) -> null;
            case BOOLEAN:
                return convertToBoolean();
            case TINYINT:
                return (dbzObj, schema) -> Byte.parseByte(dbzObj.toString());
            case SMALLINT:
                return (dbzObj, schema) -> Short.parseShort(dbzObj.toString());
            case INTEGER:
            case INTERVAL_YEAR_MONTH:
                return convertToInt();
            case BIGINT:
            case INTERVAL_DAY_TIME:
                return convertToLong();
            case DATE:
                return convertToDate();
            case TIME_WITHOUT_TIME_ZONE:
                return convertToTime();
            case TIMESTAMP_WITHOUT_TIME_ZONE:
                return convertToTimestamp(ZoneId.of("UTC"));
            case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
                return convertToLocalTimeZoneTimestamp(ZoneId.of("UTC"));
            case FLOAT:
                return convertToFloat();
            case DOUBLE:
                return convertToDouble();
            case CHAR:
            case VARCHAR:
                return convertToString();
            case BINARY:
            case VARBINARY:
                return convertToBinary();
            case DECIMAL:
                return createDecimalConverter((DecimalType) type);
            case ROW:
                return createRowConverter((RowType) type);
            case ARRAY:
            case MAP:
            case MULTISET:
            case RAW:
            default:
                throw new UnsupportedOperationException("Unsupported type: " + type);
        }
    }

    private static DeserializationRuntimeConverter convertToBoolean() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Boolean) {
                return dbzObj;
            } else if (dbzObj instanceof Byte) {
                return (byte) dbzObj == 1;
            } else if (dbzObj instanceof Short) {
                return (short) dbzObj == 1;
            } else {
                return Boolean.parseBoolean(dbzObj.toString());
            }
        };
    }

    private static DeserializationRuntimeConverter convertToInt() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Integer) {
                return dbzObj;
            } else if (dbzObj instanceof Long) {
                return ((Long) dbzObj).intValue();
            } else {
                return Integer.parseInt(dbzObj.toString());
            }
        };
    }

    private static DeserializationRuntimeConverter convertToLong() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Integer) {
                return ((Integer) dbzObj).longValue();
            } else if (dbzObj instanceof Long) {
                return dbzObj;
            } else {
                return Long.parseLong(dbzObj.toString());
            }
        };
    }

    private static DeserializationRuntimeConverter createDecimalConverter(DecimalType decimalType) {
        final int precision = decimalType.getPrecision();
        final int scale = decimalType.getScale();
        return (dbzObj, schema) -> {
            BigDecimal bigDecimal;
            if (dbzObj instanceof byte[]) {
                // decimal.handling.mode=precise
                bigDecimal = Decimal.toLogical(schema, (byte[]) dbzObj);
            } else if (dbzObj instanceof String) {
                // decimal.handling.mode=string
                bigDecimal = new BigDecimal((String) dbzObj);
            } else if (dbzObj instanceof Double) {
                // decimal.handling.mode=double
                bigDecimal = BigDecimal.valueOf((Double) dbzObj);
            } else {
                if (VariableScaleDecimal.LOGICAL_NAME.equals(schema.name())) {
                    SpecialValueDecimal decimal = VariableScaleDecimal.toLogical((Struct) dbzObj);
                    bigDecimal = decimal.getDecimalValue().orElse(BigDecimal.ZERO);
                } else {
                    // fallback to string
                    bigDecimal = new BigDecimal(dbzObj.toString());
                }
            }
            return DecimalData.fromBigDecimal(bigDecimal, precision, scale);
        };
    }

    private static DeserializationRuntimeConverter convertToDouble() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Float) {
                return ((Float) dbzObj).doubleValue();
            } else if (dbzObj instanceof Double) {
                return dbzObj;
            } else {
                return Double.parseDouble(dbzObj.toString());
            }
        };
    }

    private static DeserializationRuntimeConverter convertToFloat() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Float) {
                return dbzObj;
            } else if (dbzObj instanceof Double) {
                return ((Double) dbzObj).floatValue();
            } else {
                return Float.parseFloat(dbzObj.toString());
            }
        };
    }

    private static DeserializationRuntimeConverter convertToDate() {
        return (dbzObj, schema) -> (int) TemporalConversions.toLocalDate(dbzObj).toEpochDay();
    }

    private static DeserializationRuntimeConverter convertToTime() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Long) {
                switch (schema.name()) {
                    case MicroTime.SCHEMA_NAME:
                        return (int) ((long) dbzObj / 1000);
                    case NanoTime.SCHEMA_NAME:
                        return (int) ((long) dbzObj / 1000_000);
                }
            } else if (dbzObj instanceof Integer) {
                return dbzObj;
            }
            // get number of milliseconds of the day
            return TemporalConversions.toLocalTime(dbzObj).toSecondOfDay() * 1000;
        };
    }

    private static DeserializationRuntimeConverter convertToTimestamp(ZoneId serverTimeZone) {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof Long) {
                switch (schema.name()) {
                    case Timestamp.SCHEMA_NAME:
                        return TimestampData.fromEpochMillis((Long) dbzObj);
                    case MicroTimestamp.SCHEMA_NAME:
                        long micro = (long) dbzObj;
                        return TimestampData.fromEpochMillis(micro / 1000, (int) (micro % 1000 * 1000));
                    case NanoTimestamp.SCHEMA_NAME:
                        long nano = (long) dbzObj;
                        return TimestampData.fromEpochMillis(nano / 1000_000, (int) (nano % 1000_000));
                }
            }
            LocalDateTime localDateTime = TemporalConversions.toLocalDateTime(dbzObj, serverTimeZone);
            return TimestampData.fromLocalDateTime(localDateTime);
        };
    }

    private static DeserializationRuntimeConverter convertToLocalTimeZoneTimestamp(ZoneId serverTimeZone) {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof String) {
                String str = (String) dbzObj;
                // TIMESTAMP_LTZ type is encoded in string type
                Instant instant = Instant.parse(str);
                return TimestampData.fromLocalDateTime(LocalDateTime.ofInstant(instant, serverTimeZone));
            }
            throw new IllegalArgumentException("Unable to convert to TimestampData from unexpected value '" + dbzObj + "' of type " + dbzObj.getClass().getName());
        };
    }

    private static DeserializationRuntimeConverter convertToString() {
        return (dbzObj, schema) -> StringData.fromString(dbzObj.toString());
    }

    private static DeserializationRuntimeConverter convertToBinary() {
        return (dbzObj, schema) -> {
            if (dbzObj instanceof byte[]) {
                return dbzObj;
            } else if (dbzObj instanceof ByteBuffer) {
                ByteBuffer byteBuffer = (ByteBuffer) dbzObj;
                byte[] bytes = new byte[byteBuffer.remaining()];
                byteBuffer.get(bytes);
                return bytes;
            } else {
                throw new UnsupportedOperationException("Unsupported BYTES value type: " + dbzObj.getClass().getSimpleName());
            }
        };
    }

    private static DeserializationRuntimeConverter createRowConverter(RowType rowType) {
        final DeserializationRuntimeConverter[] fieldConverters = rowType.getFields().stream()
                .map(RowType.RowField::getType)
                .map(CustomDebeziumDeserializer::createNotNullConverter)
                .toArray(DeserializationRuntimeConverter[]::new);
        final String[] fieldNames = rowType.getFieldNames().toArray(new String[0]);

        return (dbzObj, schema) -> {
            Struct struct = (Struct) dbzObj;
            int arity = fieldNames.length;
            Row row = new Row(arity);
            for (int i = 0; i < arity; i++) {
                String fieldName = fieldNames[i];
                Field field = schema.field(fieldName);
                if (field == null) {
                    row.setField(i, null);
                } else {
                    Object fieldValue = struct.getWithoutDefault(fieldName);
                    Schema fieldSchema = schema.field(fieldName).schema();
                    Object convertedField = convertField(fieldConverters[i], fieldValue, fieldSchema);
                    row.setField(i, convertedField);
                }
            }
            return row;
        };
    }

    private static Object convertField(DeserializationRuntimeConverter fieldConverter, Object fieldValue, Schema fieldSchema) throws Exception {
        if (fieldValue == null) {
            return null;
        } else {
            return fieldConverter.convert(fieldValue, fieldSchema);
        }
    }
}

解决方案:

使用 Dinky 的 CDCSOURCE 整库同步语法进行同步。该语法和 CDAS 作用相似,可以直接自动构建一个整库入仓入湖的实时任务,并且对 source 进行了合并,不会产生额外的 MySQL 及网络压力,支持对任意 sink 的同步,如 Kafka、Doris、Hudi、JDBC 等。

使用:

  1. 依赖 jar 包:
bash 复制代码
# 将下面 Dinky 根目录下整库同步依赖包放置 $FLINK_HOME/lib 下
jar/dlink-client-base-${version}.jar
jar/dlink-common-${version}.jar
lib/dlink-client-${version}.jar
  1. 语法和参数说明:
sql 复制代码
EXECUTE CDCSOURCE jobname WITH (
  key1=val1,
  key2=val2,
  ...
)

示例说明:

  1. 实时数据合并至一个 Kafka Topic:
sql 复制代码
EXECUTE CDCSOURCE jobname WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '127.0.0.1',
  'port' = '3306',
  'username' = 'dlink',
  'password' = 'dlink',
  'checkpoint' = '3000',
  'scan.startup.mode' = 'initial',
  'parallelism' = '1',
  'table-name' = 'test.student,test.score',
  'sink.connector'='datastream-kafka',
  'sink.topic'='dlinkcdc',
  'sink.brokers'='127.0.0.1:9092'
)
  1. 实时数据同步至对应 Kafka Topic:
sql 复制代码
EXECUTE CDCSOURCE jobname WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '127.0.0.1',
  'port' = '3306',
  'username' = 'dlink',
  'password' = 'dlink',
  'checkpoint' = '3000',
  'scan.startup.mode' = 'initial',
  'parallelism' = '1',
  'table-name' = 'test.student,test.score',
  'sink.connector'='datastream-kafka',
  'sink.brokers'='127.0.0.1:9092'
)
  1. 实时数据 DataStream 入仓 Doris:
sql 复制代码
EXECUTE CDCSOURCE jobname WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '127.0.0.1',
  'port' = '3306',
  'username' = 'dlink',
  'password' = 'dlink',
  'checkpoint' = '3000',
  'scan.startup.mode' = 'initial',
  'parallelism' = '1',
  'table-name' = 'test.student,test.score',
  'sink.connector' = 'datastream-doris',
  'sink.fenodes' = '127.0.0.1:8030',
  'sink.username' = 'root',
  'sink.password' = 'dw123456',
  'sink.sink.batch.size' = '1',
  'sink.sink.max-retries' = '1',
  'sink.sink.batch.interval' = '60000',
  'sink.sink.db' = 'test',
  'sink.table.prefix' = 'ODS_',
  'sink.table.upper' = 'true',
  'sink.sink.enable-delete' = 'true'
)
  1. 实时数据 FlinkSQL 入仓 Doris:
sql 复制代码
EXECUTE CDCSOURCE jobname WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '127.0.0.1',
  'port' = '3306',
  'username' = 'dlink',
  'password' = 'dlink',
  'checkpoint' = '3000',
  'scan.startup.mode' = 'initial',
  'parallelism' = '1',
  'table-name' = 'test.student,test.score',
  'sink.connector' = 'doris',
  'sink.fenodes' = '127.0.0.1:8030',
  'sink.username' = 'root',
  'sink.password' = 'dw123456',
  'sink.sink.batch.size' = '1',
  'sink.sink.max-retries' = '1',
  'sink.sink.batch.interval' = '60000',
  'sink.sink.db' = 'test',
  'sink.table.prefix' = 'ODS_',
  'sink.table.upper' = 'true',
  'sink.table.identifier' = '${schemaName}.${tableName}',
  'sink.sink.enable-delete' = 'true'
)
  1. 实时数据入湖 Hudi:
sql 复制代码
EXECUTE CDCSOURCE demo WITH (
  'connector' = 'mysql-cdc',
  'hostname' = '127.0.0.1',
  'port' = '3306',
  'username' = 'root',
  'password' = '123456',
  'source.server-time-zone' = 'UTC',
  'checkpoint'='1000',
  'scan.startup.mode'='initial',
  'parallelism'='1',
  'database-name'='data_deal',
  'table-name'='data_deal.stu,data_deal.stu_copy1',
  'sink.connector'='hudi',
  'sink.path'='hdfs://cluster1/tmp/flink/cdcdata/${tableName}',
  'sink.hoodie.datasource.write.recordkey.field'='id',
  'sink.hoodie.parquet.max.file.size'='268435456',
  'sink.write.precombine.field'='update_time',
  'sink.write.tasks'='1',
  'sink.write.bucket_assign.tasks'='2',
  'sink.write.precombine'='true',
  'sink.compaction.async.enabled'='true',
  'sink.write.task.max.size'='1024',
  'sink.write.rate.limit'='3000',
  'sink.write.operation'='upsert',
  'sink.table.type'='COPY_ON_WRITE',
  'sink.compaction.tasks'='1',
  'sink.compaction.delta_seconds'='20',
  'sink.compaction.async.enabled'='true',
  'sink.read.streaming.skip_compaction'='true',
  'sink.compaction.delta_commits'='20',
  'sink.compaction.trigger.strategy'='num_or_time',
  'sink.compaction.max_memory'='500',
  '
end