Kafka Connect Sink能否通过SMT路由记录到新主题实现入库通知?
内置SMT能力说明
现有Kafka Connect官方内置的所有SMT均不支持将正常处理的记录路由到其他主题。内置SMT的设计定位是单条记录的字段转换、元数据修改、过滤等操作,所有输出的记录只会进入当前连接器的上下游处理链路,不存在跨主题路由的原生能力。仅死信队列配置可将处理失败的记录自动转发到指定DLQ主题,无法满足正常业务记录的主动路由需求。
自定义SMT可行性说明
不建议通过自定义SMT实现该需求,核心原因如下:
- SMT的生命周期绑定在单条记录的处理链路中,若在自定义SMT中自行嵌入生产者发送事件到控制主题,会破坏Kafka Connect的一致性保障:如果SMT发送了完成事件后,后续记录入库失败、连接器回滚偏移量,会导致客户端收到错误的完成通知,出现数据一致性问题。
- SMT本身无全局状态管理能力,无法统计同一请求批次的所有记录是否全部处理完成,只能感知单条记录的处理状态。
推荐实现方案
按不同复杂度提供三种适配方案:
- 流处理层预生成完成事件
在流处理阶段为每个请求生成唯一的request_id标识,同时统计该请求生成的总记录数。流处理完成所有记录转换后,先将转换后的业务记录发送到Sink主题,同时将request_id+总记录数的元数据发送到独立的计数主题。另外启动一个轻量消费者,监听Sink连接器的消费者组偏移量,当Sink的消费偏移量超过该request_id对应最后一条记录的偏移量时,发送完成事件到客户端订阅的控制主题。 - 扩展Sink连接器逻辑
基于现有JDBC Sink连接器的Task实现,在每批记录成功写入数据库提交偏移量后,增加批次关联request_id的完成校验逻辑,校验通过后直接发送完成事件到控制主题。该方案可以保证入库成功后才发通知,一致性最强。 - 基于数据库CDC实现通知
在业务记录入库时带上全局唯一的request_id字段,同时将request_id对应的总记录数存入单独的元数据表。通过CDC工具监听业务表的写入事件,统计每个request_id的入库条数达到总记录数时,自动发送完成事件到控制主题。该方案无需修改现有流处理和Kafka Connect的原有逻辑,侵入性最低。
内容的提问来源于stack exchange,提问作者kuro
相关产品推荐
相关产品推荐

