You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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配置和消息格式即可实现增量更新:

  1. 更换写入策略
    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
  1. 调整消息结构
    需要将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.08 19:53:41