十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

实时数据湖 flink CDC + Kafka +Doris 【企业级实战】之Source 解耦优化 【附核心源码】 03

实时数据湖 flink CDC + Kafka +Doris 【企业级实战】之Source 解耦优化 【附核心源码】 03 上一篇https://blog.csdn.net/weixin_42426485/article/details/163904839?spm1001.2014.3001.5501Source 不只是“连上数据库开始读”它决定了监听范围、初始位置、事件格式、顺序信息和后续能否可靠恢复。一、Source 模块的边界Source 层负责根据系统配置找到数据库连接根据映射配置确定需要监听的表按容量策略将表拆成多个 Source配置快照与增量启动方式将数据库变更反序列化为 CDC JSON为下游保留表身份、操作类型和顺序信息。Source 层不应该决定写哪个 Postgres 或 Doris 实例直接解析目标端 Upsert SQL自行读取 Sink 配置隐式吞掉无法解析的关键变更。二、为什么表清单来自映射配置如果在 Java 代码中写.tableList(dbo.order,dbo.order_item,dbo.customer)每次新增表都需要改代码、打包和发布。项目把表清单放到 Postgres 映射配置中由source_name source_db_type enabled选择当前系统的源表。系统配置决定连接哪个数据库 映射配置决定监听哪些表 ↓ Source Factory 构建实际 Source这样表级接入和程序发布解耦但也带来一个要求启动前必须校验映射不为空不能启动一个“正常 Running、实际零张表”的任务。三、为什么要按表分组当一个系统有几十或上百张表时把所有表放进一个 Source 会产生几个问题单 Source 初始化和快照压力集中某组异常可能影响全部表无法根据表数量和变更量拆分并行能力Job Graph 过于粗粒度。项目支持groupSize每组最多多少张表 maxGroups最多拆成多少组 sourceParallelism每个 Source 算子的并行度Factory 的核心过程可以概括为ListStringallTablesloadEnabledTables(systemName);ListListStringgroupssplitTables(allTables,groupSize,maxGroups);for(ListStringtables:groups){sources.add(buildSourceForTables(config,systemName,tables));}分组不能只看表数量。更合理的长期做法是结合变更量、表大小和业务重要性避免一张超级大表与几十张小表被机械分到同组。四、MySQL Source 的构建思路MySQL Source Factory 返回多个MySqlSourceStringMySqlCdcSourceFactoryfactorynewMySqlCdcSourceFactory(pgManager,configSourceName);ListMySqlSourceStringsourcesfactory.buildSplitSources(systemName,groupSize,maxGroups);Job 再逐组加入执行环境for(inti0;isources.size();i){DataStreamStringstreamenv.fromSource(sources.get(i),WatermarkStrategy.noWatermarks(),mysql-cdc-source-i).name(mysql-cdc-source-i).uid(sourceUid(systemName,groupTables)).setParallelism(sourceParallelism);mysqlToKafka.addKafkaSink(stream,topic,group-i);}这里最值得注意的不是循环而是 UID。只用下标group-0无法表达这一组实际包含哪些表。表清单调整后下标相同不代表算子身份相同。更安全的 UID 应包含system source type tables hash五、SQLServer Source 的差异SQLServer 不是把 MySQL 类名替换一下就结束。项目需要处理默认 Schema 通常为dbo表名可能以schema.table表达Source 使用 SQLServer Incremental Source Builder不同系统可能有特定时间字段反序列化事件顺序依赖 LSN、seqval、command_id启动前要确认 SQLServer CDC 已对库和表启用。当前 Factory 构建主干类似String[]tableArraytables.stream().map(t-t.contains(.)?t:schema.t).toArray(String[]::new);SqlServerSourceBuilderStringbuildernewSqlServerSourceBuilderString().hostname(config.getHostname()).port(config.getPort()).username(config.getUsername()).password(config.getPassword()).databaseList(config.getDatabase()).tableList(tableArray).connectTimeout(Duration.ofSeconds(10));builder.deserializer(deserializerFor(systemName));builder.startupOptions(StartupOptions.latest());returnbuilder.build();博客示例省略了敏感连接信息但代码层必须从配置加载不能散落在 Job 中。六、启动模式不是随手选一个常见启动目标包括目标含义风险Initial / Snapshot Incremental先同步存量再持续读取增量对源库快照压力大需要评估锁和一致性Latest只从任务启动后的最新位置读取会忽略已有历史数据指定位置恢复从已知 Binlog、LSN 或 Savepoint 恢复需要确保位置仍然有效项目中的某条 SQLServer Source 当前使用StartupOptions.latest()这是一项明确业务选择不应被误读为所有场景的通用答案。如果目标表为空而 Source 从 Latest 启动链路即使完全健康也永远不会自动补齐历史数据。七、反序列化为什么要在 Source 侧收敛MySQL、SQLServer 原始变更格式不同时间、二进制和删除事件也存在差异。Source 侧至少应统一输出{source:{database:...,schema:...,table:...},op:c|u|d|r,before:{},after:{},event_time:0,order:{}}这不是要求所有 Source 丢掉自身特征。SQLServer 的 LSN 等字段仍要保留只是公共字段应该稳定让 Kafka 与 Sink 不必理解多套完全不同的 JSON。八、Delete 事件为什么最容易出问题Insert 和 Update 通常可以从after获取数据Delete 往往只能从before获取主键。如果进入 Kafka 前没有正确识别 Delete或者 Redis 中缺少主键元数据就会出现Kafka Key 为空或不稳定Delete 与之前的 Update 进入不同分区下游无法生成目标端 Delete SQL任务看似成功但目标库保留了脏数据。因此主键不是性能优化信息而是 CDC 正确性的基础。九、元数据预热应在 Source 启动前完成以 MySQL 为例预热过程会从映射配置读取已启用源表查询每张表的字段和主键校验元数据不为空校验至少存在可用主键写入 Redis并记录预热数量。如果预热失败宁可阻止任务启动也不要让任务在运行时为部分表生成错误 Key。十、Source 模块的典型坑1. 表配置顺序变化导致状态错配根因通常不是 Flink 无法恢复而是 UID 没有表达真实表组身份。2. Source 已启动但没有监听表映射配置为空或enabledfalse时应直接失败而不是返回空列表。3. SQLServer 时间字段反序列化异常数据库类型、Debezium 类型和 JSON 表达可能不同需要以真实字段样例验证。4. 只验证 Insert不验证 DeleteDelete 才最能暴露主键元数据、before数据和路由问题。5. Source 分组只按表数量高频大表可能造成严重倾斜需要用运行指标反向调整分组。十一、Source 验收清单系统配置能够定位正确数据源映射配置能得到预期表清单无表时任务明确失败表分组结果和稳定 UID 可解释Insert、Update、Delete、Snapshot 事件均已验证主键和字段元数据预热成功SQLServer 顺序字段被完整保留启动模式与目标表初始化策略一致从 Checkpoint / Savepoint 恢复经过验证十二、小结Source 模块真正输出的不是一串 JSON而是一个具备明确来源、操作语义、业务 Key 和恢复位置的变更事件。下一篇继续沿数据流向下事件进入 Kafka 前为什么要 EnrichmentExactly-Once Kafka Sink 如何工作Consumer Group、Offset 和 SQLServer 保序又如何配合。
返回列表