使用UpdateOneTimestampsStrategy如何增量更新MongoDB字段?
问题描述
我需要持续更新MongoDB中_id为count的文档的value字段(初始文档为{'_id': 'count', 'value': 0}),每次增加指定数值。当前使用Kafka Connect的MongoSinkConnector,配置和代码如下:
现有Connector配置
document.id.strategy=com.mongodb.kafka.connect.sink.processor.id.strategy.ProvidedInValueStrategy writemodel.strategy=com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneTimestampsStrategy
Python生产者发送消息代码
self._aio_producer.produce( topic='mongo', value=json.dumps( { "_id":"count", "$inc":{"value":len(task['payload'].split(','))} } ) )
报错信息
Failed to put into the sink the following records: [SinkRecord{kafkaOffset=8174, timestampType=CreateTime} ConnectRecord{topic='mongo', kafkaPartition=1, key=null, keySchema=null, value={_id=count, $inc={value=1}}, valueSchema=null, timestamp=1697153679938, headers=ConnectHeaders(headers=)}] (com.mongodb.kafka.connect.sink.MongoSinkTask:244) com.mongodb.kafka.connect.sink.dlq.WriteException: v=1, code=52, message=The dollar ($) prefixed field '$inc' in '$inc' is not allowed in the context of an update's replacement document. Consider using an aggregation pipeline with $replaceWith., details={}
移除$inc后只会反复替换文档,无法实现增量更新。请问无需自定义类的情况下,是否有办法实现需求?
解决方案
无需自定义类,通过调整Connector配置和消息格式即可实现增量更新:
- 更换写入策略
UpdateOneTimestampsStrategy的逻辑是用消息内容直接替换匹配的文档,不支持MongoDB的更新操作符(如$inc)。需要将其替换为UpdateOneWritesStrategy,该策略支持解析消息中的更新操作指令。修改后的配置:
document.id.strategy=com.mongodb.kafka.connect.sink.processor.id.strategy.ProvidedInValueStrategy writemodel.strategy=com.mongodb.kafka.connect.sink.writemodel.strategy.UpdateOneWritesStrategy
- 调整消息结构
需要将MongoDB的更新操作符包装在update字段下,同时保留_id用于匹配文档(或者显式指定filter字段)。正确的消息格式示例:
self._aio_producer.produce( topic='mongo', value=json.dumps( { "_id": "count", "update": { "$inc": {"value": len(task['payload'].split(','))} } } ) )
如果需要支持文档不存在时自动初始化(替代手动插入初始文档),可以添加upsert: true:
self._aio_producer.produce( topic='mongo', value=json.dumps( { "filter": {"_id": "count"}, "update": { "$inc": {"value": len(task['payload'].split(','))} }, "upsert": True } ) )
调整后,Kafka Connect会将update字段中的内容作为MongoDB的更新命令执行,从而实现value字段的增量累加。
内容的提问来源于stack exchange,提问作者mosh_mav
相关产品推荐
相关产品推荐

