欢迎光临
我们一直在努力

实时数据湖 flink CDC + Kafka +Doris 【企业级实战】之 整体架构设计及实现 【附核心源码】 02

先看第一篇 设计思想和原始版本:https://blog.csdn.net/weixin_42426485/article/details/163904764?spm=1001.2014.3001.5501

再看整体架构

在这里插入图片描述

整条链路被拆成两个阶段:

MySQL / SQLServer

Flink CDC Source

事件标准化与业务 Key 构建

Kafka

PostgreSQL 消费链路 / Doris 消费链路 / 更多消费者

配置库、Redis 元数据和 HDFS Checkpoint 分别解决“链路怎样运行”“数据怎样解释”和“任务怎样恢复”。

架构图看起来只是多了一层 Kafka,项目边界却因此发生了根本变化:Source 只负责把数据库变化稳定地变成事件;每个目标端只负责按自己的节奏消费和落库。


一、为什么不能简单地再挂一个 Sink

在同一条 DataStream 上同时挂 PostgreSQL 和 Doris Sink,最直观:

cdcStream.sinkTo(postgresSink);
cdcStream.sinkTo(dorisSink);

但它会把三个本来可以独立变化的系统绑在一起:

数据库采集 + PostgreSQL 写入 + Doris 写入

这会带来几个问题:

  • Doris 出现延迟,背压可能沿着 Job Graph 传到 CDC Source;
  • PostgreSQL 调整批量参数,也需要重新发布整条采集链路;
  • 新增一个目标端,就要继续修改已有 Job;
  • 不同目标端无法独立保存消费进度;
  • 想重放一段历史变更,只能依赖源库日志或 Flink 状态。

真正需要拆开的不是代码文件,而是故障域、发布节奏和消费进度。

Kafka 在这里承担的就是事件边界:

上游负责“可靠地产生变化”
下游负责“按自己的方式使用变化”

二、把一个大任务拆成三类 Pipeline

项目最终形成三类逻辑链路:

source-to-kafka
kafka-to-postgres
kafka-to-doris

它们可以组合运行,也可以按生产环境的故障域拆开部署。

公共启动入口 RealtimeSyncJob 根据配置决定需要装配哪些 Pipeline:

default void run(String[] args) throws Exception {
Map<String, String> options = CdcJobSupport.parseArgs(args);
String mode = CdcJobSupport.pipelineMode(options);

RealtimeJobConfig config = RealtimeJobConfig.from(
options,
mode,
defaultSystemName(),
groupSizeEnvName(),
maxGroupsEnvName(),
consumerGroupPrefix(),
jobCode());

initializeSourceConnection(config);
KafkaConnectionManager.init();

StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();

CdcJobSupport.configureCheckpoint(env, options, ...);

if (config.isSourceToKafkaEnabled()) {
addSourceToKafkaPipeline(env, runtime);
}
if (config.isKafkaToPostgresEnabled()) {
addKafkaToPostgresPipeline(env, runtime);
}
if (config.isKafkaToDorisEnabled()) {
addKafkaToDorisPipeline(env, runtime);
}

env.execute(jobDisplayName() + ": " + config.getSystemName());
}

业务 Job 不再复制整套启动流程,只保留系统默认值和源端差异。新增业务系统时,优先增加配置,而不是再复制一个越来越长的 main()。

三、数据库变化进入 Kafka 前,先变成统一事件

MySQL 与 SQLServer 的日志格式、顺序字段和 Delete 表达方式并不相同。Kafka 如果直接存放各自的原始 JSON,下游就必须重复理解每种 Source。

因此,事件写入 Kafka 前需要完成一次标准化:

源系统与源表
业务主键 Key
操作类型 op
before / after
数据库侧顺序信息
必要的路由元数据

项目把公共处理放进 Enrichment Function:

SingleOutputStreamOperator<KafkaRecord> kafkaRecords = stream
.flatMap(new CdcKafkaRecordEnrichmentFunction(
runtime.systemName,
runtime.redisName,
runtime.redisHost,
runtime.redisPort,
runtime.redisDatabase,
runtime.redisPassword))
.name("cdc-kafka-record-enrichment-" + suffix);

Redis 中预热的主键和表结构元数据,会在这里参与业务 Key 构建。

这个 Key 不能随便生成。它至少要满足:

  • 同一张表、同一主键的事件拥有稳定 Key;
  • Insert、Update、Delete 使用同一套 Key 规则;
  • 不同表的相同主键不能相互碰撞;
  • 下游可以用它进行 keyBy、顺序判断和幂等落库。

如果 Key 不稳定,Kafka 分区内有序也救不了业务乱序。

四、Source 到 Kafka,怎样做到 Exactly Once

项目创建 Kafka Sink 时启用事务交付:

return KafkaSink.<KafkaRecord>builder()
.setBootstrapServers(bootstrapServers)
.setKafkaProducerConfig(producerProps)
.setRecordSerializer(
new CdcKafkaRecordSerializationSchema(topic))
.setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix(transactionalIdPrefix)
.build();

在正常的 Checkpoint 协调下,事件先进入 Kafka 事务;Checkpoint 成功时事务才提交。任务失败恢复后,未提交事务不会成为下游可见数据。

但 EXACTLY_ONCE 不是写上一个枚举就自动成立,它依赖:

  • Flink Checkpoint 能够持续成功;
  • transactionalIdPrefix 稳定且不会与其他算子冲突;
  • Kafka 事务超时大于可能的 Checkpoint 耗时;
  • Topic 分区和 Sink 并行度经过容量设计;
  • 下游消费者使用 read_committed 语义。

这里需要把范围说清楚:它保证的是 Source → Kafka 这一段,不会自动把 JDBC 目标端也变成端到端 Exactly Once。

五、Kafka 不是消息中转站,而是内部数据契约

Kafka 中的 CDC 事件一旦被 PostgreSQL、Doris 和其他消费者共同使用,格式就不能再随意变化。

一条可长期演进的事件至少要回答:

谁发生了变化?
哪条业务数据发生了变化?
发生了 Insert、Update 还是 Delete?
变更前后分别是什么?
事件在数据库中的先后顺序是什么?
下游应该把它路由到哪里?

因此需要把 Kafka Value 当成内部 API:

  • 新字段尽量保持向后兼容;
  • Key 生成规则不能随发布任意改变;
  • Delete 必须保留足够的主键信息;
  • 源端特有字段不能无边界地泄漏给所有下游;
  • 事件版本需要能够被消费者识别。

Kafka 真正带来的复用能力,建立在稳定事件契约之上。

六、PostgreSQL 和 Doris 为什么必须使用不同 Consumer Group

两个下游都需要收到完整 CDC 数据,因此它们必须拥有独立的 Consumer Group:

Topic: cdc-erp

Group: erp-postgres → PostgreSQL
Group: erp-doris → Doris

如果两个目标端误用了同一个 Group,它们会分摊 Topic 分区,结果不是“双写”,而是 PostgreSQL 和 Doris 各收到一部分数据。

PostgreSQL 消费链路的源码主干如下:

SingleOutputStreamOperator<KafkaRecord> kafkaStream = env.fromSource(
CdcKafkaSourceFactory.create(
config.getKafkaClusterName(),
config.getTopic(),
sinkConfig.getConsumerGroupId(),
sinkConfig.getStartupMode()),
WatermarkStrategy.noWatermarks(),
"cdc-kafka-source-" + config.getTopic())
.uid("cdc-kafka-source-" + config.getSystemName());

new KafkaToPostgres(
config.getSystemName(),
config.getConfigSourceName(),
sinkConfig.getTargetSourceName(),
config.getRedisName(),
sinkConfig.getBatchSize(),
sinkConfig.getFlushIntervalMs())
.addSink(kafkaStream, ...);

Doris 使用相同的装配模式,但拥有自己的 Group、并行度、批量大小和 Flush 周期。一个目标端扩容或恢复,不再要求另一个目标端同步操作。

七、Offset 必须交给 Checkpoint 管理

Kafka Source 关闭自动提交,并在 Checkpoint 完成时提交 Offset:

props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
props.put("commit.offsets.on.checkpoint", "true");

return KafkaSource.<KafkaRecord>builder()
.setTopics(topic)
.setGroupId(consumerGroupId)
.setStartingOffsets(toOffsetsInitializer(startupMode))
.setDeserializer(new KafkaRecordDeserializer())
.setProperties(props)
.build();

项目支持三种启动位置:

  • committed-offsets:从已提交位置继续;
  • earliest:从保留范围内最早数据开始;
  • latest:只处理启动后的新数据。

生产持续任务通常从已提交 Offset 恢复。首次上线、补数和灾难恢复则必须显式选择,不能把启动位置交给默认值碰运气。

八、Kafka 分区有序,不等于 SQLServer 业务事件一定正确

Kafka 只能保证单分区内的写入顺序。SQLServer CDC 自己还有数据库事务顺序:

start_lsn → seqval → command_id

因此,SQLServer 事件从 Kafka 读出后,落库前仍要按业务 Key 分组并判断新旧:

DataStream<KafkaRecord> ordered = stream
.keyBy(new CdcKafkaRecordKeySelector())
.process(new SqlServerKafkaRecordOrderProcessFunction());

这一层会丢弃已经处理过的旧事件,顺序状态则随 Checkpoint 恢复。

这也说明了一个容易误解的点:Kafka 提供的是传输层顺序能力,数据库语义层的顺序仍然要由 CDC 链路自己解释。

九、Kafka 到 JDBC 的一致性边界在哪里

PostgreSQL 和 Doris 当前采用的不是两阶段提交 Sink,而是:

At-least-once + 目标端幂等

处理顺序大致如下:

Kafka 事件进入 Buffer

按目标表批量 Upsert / Delete

Checkpoint 前强制 Flush

数据库事务 Commit

Checkpoint 完成并提交 Offset

如果数据库已经 Commit,但 Checkpoint 在完成前失败,任务恢复后可能再次消费这批事件。

所以端到端结果正确依赖目标表主键与 Upsert:重复 Insert/Update 要收敛到同一条记录,Delete 要按同一业务主键执行。

这不是严格意义上的 JDBC Exactly Once,而是通过重放与幂等获得最终一致。

十、配置和元数据也要跟着架构分层

项目中有三类不同的配置事实:

配置载体管理内容作用
Pipeline YAML Source、Topic、Consumer Group、Sink 开关、目标实例 描述链路如何连接
Resource YAML 并行度、表分组、重启和资源参数 描述链路如何运行
PostgreSQL 映射表 源表、目标表和字段路由 描述数据写到哪里

Redis 保存主键、字段和类型等运行期元数据,用于:

  • 构建 Kafka 业务 Key;
  • 解释 Delete;
  • 生成目标端 Upsert;
  • 完成字段类型绑定;
  • 让多个消费者复用同一份表结构事实。

配置库决定“走哪条路”,元数据决定“这条数据究竟是什么”。

十一、这一层 Kafka 到底换来了什么

1. 采集和写入可以独立发布

修改 Doris Sink 不需要重启数据库 CDC Source;新增消费者也不必改动已有采集链路。

2. 下游故障不再立即阻断上游

在 Kafka 保留时间和容量允许的范围内,下游可以暂时停止,恢复后从自己的 Offset 继续消费。

3. 每个目标端拥有独立进度

PostgreSQL 可以追到最新,Doris 可以暂时落后;两者不再共享一个处理位置。

4. 数据可以重放

在 Topic 保留范围内,可以通过新 Consumer Group 或重置 Offset 进行补偿、回归验证和新目标初始化。

5. Source 接入方式得到统一

MySQL 和 SQLServer 在进入 Kafka 前收敛为统一 CDC 事件,下游不必为每种数据库复制一套完整链路。

这些收益不是“消息队列性能更高”这么简单,而是系统边界变清楚了。

十二、链路跑起来后,新的坑也出现了

Kafka 解决了耦合问题,也带来了自己的工程成本:

  • Topic 创建成功后,元数据传播并不一定立即完成;
  • Transactional ID 不稳定或冲突,会影响事务恢复;
  • Consumer Group 写错,可能让两个目标端意外分摊数据;
  • Partition 数、Source 并行度和 Sink 并行度不匹配,会产生热点;
  • Kafka Key 规则变化,可能破坏同主键事件顺序;
  • Redis 主键元数据缺失,会影响 Delete、路由和幂等;
  • 事件格式不兼容,会同时影响多个下游;
  • Kafka Lag 变大后,需要判断是消费能力不足、目标库变慢,还是某个分区倾斜;
  • Offset、Checkpoint 和 JDBC Commit 之间仍存在一致性边界。

于是,整体架构虽然清晰了,真正决定生产稳定性的细节才刚刚开始。

后面的文章会逐个拆开这些模块:Source 表发现与分组、Kafka Key 与事务、配置和 Redis 元数据、Checkpoint 与恢复、JDBC Sink 抽象、典型踩坑,以及最后的性能和运维体系。

赞(0)
未经允许不得转载:171主机测评 » 实时数据湖 flink CDC + Kafka +Doris 【企业级实战】之 整体架构设计及实现 【附核心源码】 02
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址