如何通过Kafka Connectors按需查询MongoDB等数据库数据至Kafka主题?
使用Kafka Connectors从MongoDB按需查询指定数据到Kafka主题
核心结论
完全可以通过Kafka Connect的MongoDB Source Connector实现按需查询指定数据并写入Kafka主题,无需编写额外业务代码,仅需通过配置即可完成。
具体实现步骤
1. 配置MongoDB Source Connector
创建一个Connector配置文件(例如mongodb-mahindra-car-source.properties),关键配置如下:
name=mahindra-car-query-source connector.class=com.mongodb.kafka.connect.MongoSourceConnector tasks.max=1 connection.uri=mongodb://你的MongoDB地址:27017/ database=你的数据库名 collection=汽车集合名 # 指定目标Kafka主题 topic=mahindra-cars-topic # 用MongoDB聚合管道过滤Mahindra品牌数据 pipeline=[{"$match": {"Brand": "Mahindra"}}] # 设置为单次运行模式,完成查询后自动停止 batch.size=0 poll.max.batch.size=1000 poll.await.time.ms=0
- pipeline:通过MongoDB的
$match聚合操作精准筛选Brand为Mahindra的文档,支持扩展更复杂的查询逻辑。 - batch.size=0与poll.await.time.ms=0:让Connector仅执行一次查询任务,获取所有符合条件的数据后立即停止,实现"按需触发"的效果。
2. 提交Kafka Connect任务
通过Kafka Connect的REST API提交配置:
curl -X POST -H "Content-Type: application/json" --data @mongodb-mahindra-car-source.properties http://你的Connect地址:8083/connectors
3. 验证结果
使用Kafka消费者命令查看目标主题的数据:
kafka-console-consumer.sh --bootstrap-server 你的Kafka地址:9092 --topic mahindra-cars-topic --from-beginning
注意事项
- 确保MongoDB Kafka Connector已安装到Kafka Connect的插件目录中。
- 若需重复执行查询,重新提交相同配置任务即可。
- 如需持续监听新增的Mahindra品牌数据,可移除
batch.size=0等单次运行配置,让Connector持续同步符合条件的新数据。
内容的提问来源于stack exchange,提问作者Rathore
相关产品推荐
相关产品推荐

