是否可通过KsqlDB拆分Kafka单Topic消息写入MySQL多表
结论
完全可行。你描述的链路可以正常落地,用KsqlDB做流拆分、关联补全自增ID,再对接JDBC Sink写入MySQL的方案,是这类场景的常规实现方式。
具体实现流程
- 首先给原始CSV主题注册KsqlDB流
先把Topic1映射为KsqlDB流,明确字段定义,注意序列化格式要和SFTP Connector输出的格式保持一致:CREATE STREAM s_raw_csv ( Col1 VARCHAR, Col2 VARCHAR, Col3 VARCHAR, Col4 VARCHAR ) WITH ( KAFKA_TOPIC = 'Topic1', VALUE_FORMAT = 'CSV', -- 如果SFTP Connector输出为JSON/AVRO对应修改即可 PARTITIONS = 6 -- 和Topic1实际分区数保持一致 ); - 注册CDC同步的MySQL表为KsqlDB表
把两个MySQL表通过CDC同步到Kafka的changelog主题注册为KsqlDB的只读表,用于后续关联查询自增ID:-- 注册Table2的CDC主题,主键为MySQL自增Id CREATE TABLE t_cdc_table2 ( Id BIGINT PRIMARY KEY, Col3 VARCHAR, Col4 VARCHAR ) WITH ( KAFKA_TOPIC = 'cdc_tbl2', -- 替换为你实际CDC同步Table2的主题名 VALUE_FORMAT = 'JSON' -- 和CDC输出格式保持一致 ); -- 注册Table1的CDC主题,主键为MySQL自增Id CREATE TABLE t_cdc_table1 ( Id BIGINT PRIMARY KEY, Col1 VARCHAR, Col2 VARCHAR, Tbl2Id BIGINT ) WITH ( KAFKA_TOPIC = 'cdc_tbl1', -- 替换为你实际CDC同步Table1的主题名 VALUE_FORMAT = 'JSON' ); - 拆分生成Table2对应的Sink主题
从原始流中抽取Col3、Col4字段,左关联CDC的Table2表补全已存在的自增ID,输出到独立主题供Table2的Sink消费:CREATE STREAM s_sink_tbl2 WITH ( KAFKA_TOPIC = 'sink_topic_table2', VALUE_FORMAT = 'JSON', PARTITIONS = 6 ) AS SELECT COALESCE(t2.Id, NULL) AS Id, -- 已有记录取原ID,新记录传NULL由MySQL自增生成 s.Col3 AS Col3, s.Col4 AS Col4 FROM s_raw_csv s LEFT JOIN t_cdc_table2 t2 -- 关联条件替换为你Table2实际的业务唯一键匹配规则,比如Col3+Col4组合唯一 ON s.Col3 = t2.Col3 AND s.Col4 = t2.Col4 PARTITION BY COALESCE(t2.Id, s.Col3); -- 按主键分区保证同记录消息有序 - 拆分生成Table1对应的Sink主题
从原始流中抽取Col1、Col2字段,关联两个CDC表分别补全Table1自身的自增ID、关联的Table2 ID,输出到独立主题供Table1的Sink消费:CREATE STREAM s_sink_tbl1 WITH ( KAFKA_TOPIC = 'sink_topic_table1', VALUE_FORMAT = 'JSON', PARTITIONS = 6 ) AS SELECT COALESCE(t1.Id, NULL) AS Id, -- 已有记录取原ID,新记录传NULL由MySQL自增生成 s.Col1 AS Col1, s.Col2 AS Col2, t2.Id AS Tbl2Id -- 关联拿到对应Table2的自增ID FROM s_raw_csv s LEFT JOIN t_cdc_table1 t1 ON s.Col1 = t1.Col1 -- 按要求通过Col1匹配Table1已有记录 LEFT JOIN t_cdc_table2 t2 ON s.Col3 = t2.Col3 AND s.Col4 = t2.Col4 PARTITION BY COALESCE(t1.Id, s.Col1);
配置注意事项
- 数据库层面必须给业务键加唯一索引:Table1的Col1字段要加唯一索引,Table2的业务匹配字段(比如Col3+Col4)也要加唯一索引,避免关联时匹配到多条记录,同时也能保证JDBC Sink的upsert逻辑正常执行。
- JDBC Sink配置:两个表对应的Sink都要设置
insert.mode=upsert,pk.mode=record_value,pk.fields=Id,同时关闭自动建表、自动改表配置,写入新记录时Id传NULL不会触发报错,MySQL会自动生成自增ID。 - 新记录写入时序问题:第一次写入全新的Col1记录时,因为Table2的自增ID需要等MySQL插入完成、CDC同步回Kafka后才能被关联到,会出现第一次写入Table1时Tbl2Id为空的情况。可以给两个流的JOIN配置5-10秒的滚动窗口,等Table2的ID回传后再输出Table1的写入消息即可规避。
- 分区配置:所有输出的Sink主题分区数要和原始Topic1保持一致,且必须按表的主键分区,避免同一条业务记录的消息乱序到达Sink,触发更新先于插入的异常。
内容的提问来源于stack exchange,提问作者Umesh N
相关产品推荐
相关产品推荐

