FlinkCDC实战进阶指南
Flink CDC 同步 MySQL 到 Doris 之 Schema 变更
Flink CDC DataStream 双流 Join

背景
在实际生产环境中,表的 Schema 信息经常会被修改,而且整库同步时,随着业务的增长,MySQL 的压力也会越来越大。重启任务和集群、重新消费显然是不合理的,因此在做 Flink CDC 时需要兼顾并解决 Schema 变更、增删表的问题。
解决方案
方案 1: Flink CDC 读取时开启 Schema Change,并且 Doris 建表时也开启 Light Schema Change
-
Flink CDC 参数
java.scanNewlyAddedTableEnabled(true) // 启用扫描新添加表的功能 .includeSchemaChanges(true) -
Doris 建表时指定
sqllight_schema_change=true -
程序读取 MySQL 中获取需要同步的表,以字段
member_id,table字段存储 Doris 中表 A。 -
脚本读取 Doris 表 A 数据,获取 MySQL 中的 Schema,通过转换,获取 Doris 建表语句,连接 Doris 执行语句。
-
取消 Flink 任务,并重新启动 Flink 任务(重启只适合添加新库,新表不用重启)。
-
每次重启连接 Doris 表 A,获取 database,组装
databaseList,tableList,tableList使用正则,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 解析和转化
-
实现自己的
DebeziumDeserializationSchema,需要实现deserialize、getProducedType两个函数。deserialize实现转换数据的逻辑。getProducedType定义返回的类型,这里返回两个参数,第一个 Boolean 类型的参数表示数据是upsert或是delete,第二个参数返回转换后的 JSON string,这里的 JSON 将会包含 Schema 变更后的 Column 与对应的 Value。
-
如果启动时设置的
.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");
}
}

end
