改版通知

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

FlinkCDC实战进阶指南

ckckck2025年1月10日46 浏览

Flink CDC 同步 MySQL 到 Doris 之 Schema 变更

Flink CDC 同步 Doris

背景

在实际生产环境中,表的 Schema 信息经常会被修改,而且整库同步时,随着业务的增长,MySQL 的压力也会越来越大。重启任务和集群、重新消费显然是不合理的,因此在做 Flink CDC 时需要兼顾并解决 Schema 变更、增删表的问题。

解决方案

  • Flink CDC 参数

    java 复制代码
    .scanNewlyAddedTableEnabled(true)  // 启用扫描新添加表的功能
    .includeSchemaChanges(true)
  • Doris 建表时指定

    sql 复制代码
    light_schema_change=true
  • 程序读取 MySQL 中获取需要同步的表,以字段 member_id, table 字段存储 Doris 中表 A。

  • 脚本读取 Doris 表 A 数据,获取 MySQL 中的 Schema,通过转换,获取 Doris 建表语句,连接 Doris 执行语句。

  • 取消 Flink 任务,并重新启动 Flink 任务(重启只适合添加新库,新表不用重启)。

  • 每次重启连接 Doris 表 A,获取 database,组装 databaseList, tableListtableList 使用正则,database1.*, database2.*,对库内所有表进行监听,这样可以达到 MySQL 添加新表时将新表加入同步队列。

  • Doris 目前已经支持 Schema 变更了,只不过 CDC DDL 变更获取需要自行实现,然后通过 JDBC 的方式连接 Doris 去执行 DDL SQL,因为 SQL 有点差异,需要转换才能执行,结合 MySQL 新表,可以在 DDL 获取 CREATE 对 Doris 进行建表。

  • 在将数据导入到 Doris 时,速度导入过快都会出现导入失败,-235 错误,可以使用控制读取 Binlog 数量 + Window 聚合去批量导入。

  • 如需要导入表 B 的数据有 {"id":1,"name":"小明"}, {"id":2,"name":"小红"},如果执行两次 PUT 显然是不合理的,可以使用 jsonArr 的方式 [{"id":1,"name":"小明"},{"id":2,"name":"小红"}] 一次导入。

方案 2: 自定义序列化方式,进行 Schema 解析和转化

  1. 实现自己的 DebeziumDeserializationSchema,需要实现 deserializegetProducedType 两个函数。

    • deserialize 实现转换数据的逻辑。
    • getProducedType 定义返回的类型,这里返回两个参数,第一个 Boolean 类型的参数表示数据是 upsert 或是 delete,第二个参数返回转换后的 JSON string,这里的 JSON 将会包含 Schema 变更后的 Column 与对应的 Value。
  2. 如果启动时设置的 .serverTimeZone("Asia/Shanghai") 并没有生效,查源码可以发现,底层的 Debezium 并没有实现 serverTimeZone 的配置,相应的转换是在 RowDataDebeziumDeserializeSchema 内实现的。

java 复制代码
interface DeserializationRuntimeConverter extends Serializable {
    Object convert(Object dbzObj, Schema schema);
}

import com.ververica.cdc.debezium.DebeziumDeserializationSchema;
import com.ververica.cdc.debezium.utils.TemporalConversions;
import io.debezium.data.Envelope;
import io.debezium.time.Date;
import io.debezium.time.MicroTimestamp;
import io.debezium.time.NanoTimestamp;
import io.debezium.time.Timestamp;
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.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.flink.table.data.TimestampData;
import org.apache.flink.util.Collector;
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.time.ZoneOffset;
import java.time.format.DateTimeFormatter;
import java.util.Map;
import java.util.stream.Collectors;

public class JsonStringDebeziumDeserializationSchema implements DebeziumDeserializationSchema {
    int zoneOffset;

    @Override
    public void deserialize(SourceRecord record, Collector out) throws Exception {
        Envelope.Operation op = Envelope.operationFor(record);
        Struct value = (Struct) record.value();
        Schema valueSchema = record.valueSchema();
        if (op == Envelope.Operation.CREATE || op == Envelope.Operation.READ) {
            String insert = extractAfterRow(value, valueSchema);
            out.collect(new Tuple2<>(true, insert));
        } else if (op == Envelope.Operation.DELETE) {
            String delete = extractBeforeRow(value, valueSchema);
            out.collect(new Tuple2<>(false, delete));
        } else {
            String after = extractAfterRow(value, valueSchema);
            out.collect(new Tuple2<>(true, after));
        }
    }

    public JsonStringDebeziumDeserializationSchema() {
        // 实现一个用于转换时间的 Converter
        this.runtimeConverter = (dbzObj, schema) -> {
            if (schema.name() != null) {
                switch (schema.name()) {
                    case Timestamp.SCHEMA_NAME:
                        return TimestampData.fromEpochMillis((Long) dbzObj).toLocalDateTime().atOffset(ZoneOffset.ofHours(zoneOffset)).format(DateTimeFormatter.ISO_OFFSET_DATE_TIME);
                    case MicroTimestamp.SCHEMA_NAME:
                        long micro = (long) dbzObj;
                        return TimestampData.fromEpochMillis(micro / 1000, (int) (micro % 1000 * 1000)).toLocalDateTime().atOffset(ZoneOffset.ofHours(zoneOffset)).format(DateTimeFormatter.ISO_OFFSET_DATE_TIME);
                    case NanoTimestamp.SCHEMA_NAME:
                        long nano = (long) dbzObj;
                        return TimestampData.fromEpochMillis(nano / 1000_000, (int) (nano % 1000_000)).toLocalDateTime().atOffset(ZoneOffset.ofHours(zoneOffset)).format(DateTimeFormatter.ISO_OFFSET_DATE_TIME);
                    case Date.SCHEMA_NAME:
                        return TemporalConversions.toLocalDate(dbzObj).format(DateTimeFormatter.ISO_LOCAL_DATE);
                }
            }
            return dbzObj;
        };
    }

    private final DeserializationRuntimeConverter runtimeConverter;

    private Map<String, Object> getRowMap(Struct after) {
        // 转换时使用对应的转换器
        return after.schema().fields().stream().collect(Collectors.toMap(Field::name, f -> after.get(f)));
    }

    private String extractAfterRow(Struct value, Schema valueSchema) throws Exception {
        Struct after = value.getStruct(Envelope.FieldName.AFTER);
        Map<String, Object> rowMap = getRowMap(after);
        ObjectMapper objectMapper = new ObjectMapper();
        return objectMapper.writeValueAsString(rowMap);
    }

    private String extractBeforeRow(Struct value, Schema valueSchema) throws Exception {
        Struct after = value.getStruct(Envelope.FieldName.BEFORE);
        Map<String, Object> rowMap = getRowMap(after);
        ObjectMapper objectMapper = new ObjectMapper();
        return objectMapper.writeValueAsString(rowMap);
    }

    @Override
    public TypeInformation getProducedType() {
        return TypeInformation.of(new TypeHint<Tuple2<Boolean, String>>() {});
    }
}
Flink CDC 双流 Join
Flink CDC 常见的三种 Join 场景

MySQL 数据准备

sql 复制代码
# 创建数据库 flinkcdc_etl_test
CREATE DATABASE flinkcdc_etl_test;

# 使用数据库 flinkcdc_etl_test
USE flinkcdc_etl_test;

# 创建教师表
DROP TABLE IF EXISTS `teacher`;
CREATE TABLE `teacher` (
  `t_id` varchar(3) NOT NULL COMMENT '主键',
  `t_name` varchar(10) NOT NULL COMMENT '教师名称',
  PRIMARY KEY (`t_id`) USING BTREE
) COMMENT = '教师表';

# 插入教师表数据
INSERT INTO `teacher` VALUES ('001', '张三');
INSERT INTO `teacher` VALUES ('002', '李四');
INSERT INTO `teacher` VALUES ('003', '王五');

# 创建课程表
DROP TABLE IF EXISTS `course`;
CREATE TABLE `course` (
  `c_id` varchar(3) NOT NULL COMMENT '主键',
  `c_name` varchar(20) NOT NULL COMMENT '课程名称',
  `c_tid` varchar(3) NOT NULL COMMENT '教师表主键',
  PRIMARY KEY (`c_id`) USING BTREE,
  INDEX `c_tid`(`c_tid`) USING BTREE
) COMMENT = '课程表';

# 插入课程表数据
INSERT INTO `course` VALUES ('1', '语文', '001');
INSERT INTO `course` VALUES ('2', '数学', '002');

代码实现

1. 基础类

java 复制代码
/**
 * CDC 中 op 类型
 */
public enum OpEnum {
    /**
     * 新增
     */
    CREATE("c", "create", "新增"),
    /**
     * 修改
     */
    UPDATA("u", "update", "更新"),
    /**
     * 删除
     */
    DELETE("d", "delete", "删除"),
    /**
     * 读
     */
    READ("r", "read", "首次读");

    /**
     * 字典码
     */
    private String dictCode;

    /**
     * 字典码翻译值
     */
    private String dictValue;

    /**
     * 字典码描述
     */
    private String description;

    OpEnum(String dictCode, String dictValue, String description) {
        this.dictCode = dictCode;
        this.dictValue = dictValue;
        this.description = description;
    }

    public String getDictCode() {
        return dictCode;
    }

    public String getDictValue() {
        return dictValue;
    }

    public String getDescription() {
        return description;
    }
}

2. 工具类

java 复制代码
import com.alibaba.fastjson.JSONObject;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

import java.time.Duration;

/**
 * 转换工具类
 */
public class TransformUtil {

    /**
     * 格式化抽取数据格式
     * 去除 before、after、source 等冗余内容
     *
     * @param extractData 抽取的数据
     * @return
     */
    public static JSONObject formatResult(String extractData) {
        JSONObject formatDataObj = new JSONObject();
        JSONObject rawDataObj = JSONObject.parseObject(extractData);
        formatDataObj.putAll(rawDataObj);
        formatDataObj.remove("before");
        formatDataObj.remove("after");
        formatDataObj.remove("source");
        String op = rawDataObj.getString("op");
        if (OpEnum.DELETE.getDictCode().equals(op)) {
            // 新增取 before 结构体数据
            formatDataObj.putAll(rawDataObj.getJSONObject("before"));
        } else {
            // 其余取 after 结构体数据
            formatDataObj.putAll(rawDataObj.getJSONObject("after"));
        }
        return formatDataObj;
    }

    static MySqlSource<String> getStringMySqlSource(String dbNasme, String tableName) {
        MySqlSource<String> teacherSouce = MySqlSource.<String>builder()
                .hostname("192.168.18.101")
                .port(3306)
                .username("root")
                .password("123456")
                .databaseList(dbNasme)
                .tableList(dbNasme + "." + tableName)
                .startupOptions(StartupOptions.initial())
                .deserializer(new JsonDebeziumDeserializationSchema())
                .serverTimeZone("Asia/Shanghai")
                .build();
        return teacherSouce;
    }

    static DataStreamSource<String> getStringDataStreamSource(MySqlSource<String> dataStream, StreamExecutionEnvironment env) {
        DataStreamSource<String> mysqlDataStreamSource = env.fromSource(
                dataStream,
                WatermarkStrategy.<String>forBoundedOutOfOrderness(Duration.ofSeconds(1L)).withTimestampAssigner(
                        new SerializableTimestampAssigner<String>() {
                            @Override
                            public long extractTimestamp(String extractData, long l) {
                                return JSONObject.parseObject(extractData).getLong("ts_ms");
                            }
                        }
                ),
                "DataStreamWithWatermark Source"
        );
        return mysqlDataStreamSource;
    }
}

3. 基于事件时间的窗口内关联

java 复制代码
import com.alibaba.fastjson.JSONObject;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.api.common.functions.JoinFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import static org.apache.doris.flink.test.TransformUtil.*;

/**
 * 基于事件时间的 window inner join 把教师表和课程表进行联结
 * <p>
 * 只有两者数据流关联到数据,才会进行打印
 */
public class WindowInnerJoinByEventTimeTest {

    public static void main(String[] args) throws Exception {
        // 1. 创建 Flink-MySQL-CDC 的 Source
        MySqlSource<String> teacherSouce = getStringMySqlSource("flinkcdc_etl_test", "teacher");
        MySqlSource<String> courseSouce = getStringMySqlSource("flinkcdc_etl_test", "course");

        // 2. 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // 3. 提取 waterMark
        DataStreamSource<String> teacherDataStreamSource = getStringDataStreamSource(teacherSouce, env);
        DataStreamSource<String> courseDataStreamSource = getStringDataStreamSource(courseSouce, env);

        // 4. 转换为指定格式
        DataStream<JSONObject> teacherDataStream = teacherDataStreamSource.map(rawData -> formatResult(rawData));
        DataStream<JSONObject> courseDataStream = courseDataStreamSource.map(rawData -> formatResult(rawData));

        // 5. 窗口联结(教师流和课程表)打印输出
        windowInnerJoinAndPrint(teacherDataStream, courseDataStream);

        // 6. 执行任务
        env.execute("WindowInnerJoinByEventTimeTest Job");
    }

    /**
     * 窗口联结并打印输出
     * 只支持 inner join,即窗口内联关联到的才会下发,关联不到的则直接丢掉。
     * 如果想实现 Window 上的 outer join,需要使用 coGroup 算子
     *
     * @param teacherDataStream 教师数据流
     * @param courseDataStream  课程数据流
     */
    private static void windowInnerJoinAndPrint(DataStream<JSONObject> teacherDataStream,
                                                DataStream<JSONObject> courseDataStream) {
        DataStream<JSONObject> teacherCourseDataStream = teacherDataStream
                .join(courseDataStream)
                .where(teacher -> teacher.getString("t_id"))
                .equalTo(couse -> couse.getString("c_tid"))
                .window(TumblingEventTimeWindows.of(Time.seconds(10L)))
                .apply(
                        new JoinFunction<JSONObject, JSONObject, JSONObject>() {
                            @Override
                            public JSONObject join(JSONObject jsonObject,
                                                   JSONObject jsonObject2) {
                                // 拼接
                                jsonObject.putAll(jsonObject2);
                                return jsonObject;
                            }
                        }
                );
        teacherCourseDataStream.print("Window Inner Join By Event Time");
    }
}

4. 基于事件时间的窗口外关联

java 复制代码
import com.alibaba.fastjson.JSONObject;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.api.common.functions.CoGroupFunction;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.util.Collector;
import static org.apache.doris.flink.test.TransformUtil.*;

/**
 * 根据 event time(事件时间)进行 window outer join(窗口外关联)
 * 把教师表和课程表进行窗口外联联结,关联不到的数据也会下发
 */
public class WindowOuterJoinByEventTimeTest {

    public static void main(String[] args) throws Exception {
        // 1. 创建 Flink-MySQL-CDC 的 Source
        MySqlSource<String> teacherSouce = getStringMySqlSource("flinkcdc_etl_test", "teacher");
        MySqlSource<String> courseSouce = getStringMySqlSource("flinkcdc_etl_test", "course");

        // 2. 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // 3. 提取 waterMark
        DataStreamSource<String> teacherDataStreamSource = getStringDataStreamSource(teacherSouce, env);
        DataStreamSource<String> courseDataStreamSource = getStringDataStreamSource(courseSouce, env);

        // 4. 转换为指定格式
        DataStream<JSONObject> teacherDataStream = teacherDataStreamSource.map(rawData -> formatResult(rawData));
        DataStream<JSONObject> courseDataStream = courseDataStreamSource.map(rawData -> formatResult(rawData));

        // 5. 窗口联结(教师流和课程表)打印输出
        windowOuterJoinAndPrint(teacherDataStream, courseDataStream);

        // 6. 执行任务
        env.execute("WindowOuterJoinByEventTimeTest Job");
    }

    /**
     * 窗口外联并打印输出
     * Window 上的 outer join,使用 coGroup 算子,关联不到的数据也会下发
     *
     * @param teacherDataStream 教师数据流
     * @param courseDataStream  课程数据流
     */
    private static void windowOuterJoinAndPrint(DataStream<JSONObject> teacherDataStream,
                                                DataStream<JSONObject> courseDataStream) {
        DataStream<JSONObject> teacherCourseDataStream = teacherDataStream
                .coGroup(courseDataStream)
                .where(teacher -> teacher.getString("t_id"))
                .equalTo(course -> course.getString("c_tid"))
                .window(TumblingEventTimeWindows.of(Time.seconds(10L)))
                .apply(
                        new CoGroupFunction<JSONObject, JSONObject, JSONObject>() {
                            @Override
                            public void coGroup(Iterable<JSONObject> iterable,
                                                Iterable<JSONObject> iterable1,
                                                Collector<JSONObject> collector) {
                                JSONObject result = new JSONObject();
                                for (JSONObject jsonObject : iterable) {
                                    result.putAll(jsonObject);
                                }
                                for (JSONObject jsonObject : iterable1) {
                                    result.putAll(jsonObject);
                                }
                                collector.collect(result);
                            }
                        }
                );
        teacherCourseDataStream.print("Window Outer Join By Event Time");
    }
}

5. 间隔联结

间隔联结
java 复制代码
import com.alibaba.fastjson.JSONObject;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.util.Collector;
import static org.apache.doris.flink.test.TransformUtil.*;

/**
 * interval join(间隔联结)把教师表和课程表进行联结
 * 间隔联结只支持事件时间,不支持处理时间
 */
public class InteralJoinByEventTimeTest {

    public static void main(String[] args) throws Exception {
        // 1. 创建 Flink-MySQL-CDC 的 Source
        MySqlSource<String> teacherSouce = getStringMySqlSource("flinkcdc_etl_test", "teacher");
        MySqlSource<String> courseSouce = getStringMySqlSource("flinkcdc_etl_test", "course");

        // 2. 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        // 3. 提取 waterMark
        DataStreamSource<String> teacherDataStreamSource = getStringDataStreamSource(teacherSouce, env);
        DataStreamSource<String> courseDataStreamSource = getStringDataStreamSource(courseSouce, env);

        // 4. 转换为指定格式
        DataStream<JSONObject> teacherDataStream = teacherDataStreamSource.map(rawData -> formatResult(rawData));
        DataStream<JSONObject> courseDataStream = courseDataStreamSource.map(rawData -> formatResult(rawData));

        // 3. 间隔联结(教师流和课程表)打印输出
        intervalJoinAndPrint(teacherDataStream, courseDataStream);

        // 4. 执行任务
        env.execute("TeacherJoinCourseTest Job");
    }

    /**
     * 间隔联结并打印输出
     *
     * @param teacherDataStream 教师数据流
     * @param courseDataStream  课程数据流
     */
    private static void intervalJoinAndPrint(DataStream<JSONObject> teacherDataStream,
                                             DataStream<JSONObject> courseDataStream) {
        DataStream<JSONObject> teacherCourseDataStream = teacherDataStream
                .keyBy(teacher -> teacher.getString("t_id"))
                .intervalJoin(
                        courseDataStream.keyBy(course -> course.getString("c_tid"))
                )
                .between(
                        Time.seconds(-5),
                        Time.seconds(5)
                )
                .process(
                        new ProcessJoinFunction<JSONObject, JSONObject, JSONObject>() {
                            @Override
                            public void processElement(JSONObject left, JSONObject right,
                                                       Context ctx, Collector<JSONObject> out) {
                                left.putAll(right);
                                out.collect(left);
                            }
                        }
                );
        teacherCourseDataStream.print("Interval Join By Event Time");
    }
}
Flink CDC 双流 Join

end