Kafka Streams是否支持读取非Kafka主题数据源进行流处理?
问题场景
我有一个名为smscb-router的应用,架构如下所示,核心处理逻辑为:
- 从遗留短信系统(sms)中读取数据
- 根据数据内容标识的回调类型,将数据路由写入对应的输出主题,例如
billing-n-cdr、dr-cdr等
我判断Kafka Streams API非常适配这个场景,它提供的map算子可以很便捷地实现内容映射、校验逻辑。但我存在一个疑问:目前我在网络博客中看到的所有流处理应用示例,均是从源Kafka主题读取数据,经处理后写入其他目标主题。Kafka Streams是否支持从非Kafka主题的数据源读取源数据?例如对接Redis存储、RabbitMQ这类消息队列作为输入源?

解答
原生Kafka Streams 不支持直接对接非Kafka数据源作为输入源。
Kafka Streams从设计之初就是深度绑定Kafka生态的流处理类库,它的消费位点管理、状态存储、容错恢复、处理语义保障(比如恰好一次语义)全部依托Kafka的分区机制、消费组协议实现,原生API仅提供了从Kafka主题构建KStream/KTable的入口,没有开放对接第三方系统的原生Source接口。
针对你的smscb-router场景,生产环境有两种成熟的落地方案:
- 前置轻量同步任务
单独实现一个逻辑极简的同步服务,负责从遗留短信系统、RabbitMQ、Redis等非Kafka源拉取数据,做基础格式转换后写入指定的Kafka输入主题,后续所有路由、校验逻辑全部通过Kafka Streams实现,从该输入主题消费数据即可。这个方案可以完全复用Kafka Streams自带的高可用、状态管理、语义保障能力,核心路由逻辑的开发维护成本最低,是最常用的实现方式。 - 基于Kafka Connect做异构源同步
如果你对接的是RabbitMQ、Redis这类通用数据源,不需要自行开发同步服务,可以直接使用对应数据源官方/社区提供的Kafka Connect插件,通过配置化的方式完成非Kafka源到Kafka主题的数据同步,后续流处理逻辑仍由Kafka Streams承担。
注意不要尝试在Kafka Streams的处理拓扑中硬编码拉取第三方数据源的逻辑,这种做法会破坏Kafka Streams原生的位点管理、容错恢复机制,极易引发数据重复消费、数据丢失、任务恢复异常等问题,不推荐在生产环境使用。
内容的提问来源于stack exchange,提问作者Mopparthy Ravindranath
相关产品推荐
相关产品推荐

