首页/新闻资讯/正文详情

Apache Pulsar MongoDB Sink Connector 完全指南:从配置到源码原理

发布时间:2026/9/24 20:25:26 来源:云帆数科 栏目:资讯中心
Apache Pulsar MongoDB Sink Connector 完全指南:从配置到源码原理
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读MongoDB sink connector 是 Apache Pulsar 内置的 IO 连接器之一它的职责是从 Pulsar topic 中持续拉取消息并将消息以 JSON 文档的形式批量写入 MongoDB 集合collection。本文以官方文档 io-mongo-sink.md 为骨架结合仓库中 MongoSink.java 与 MongoConfig.java 的真实实现逐一讲解每个配置项的含义与默认值、JSON/YAML 配置文件的编写方法、消息批量写入与确认ack/fail的底层机制并给出基于pulsar-admin的本地运行与部署示例。读完本文你将能够独立编写 MongoDB sink 配置、理解其批处理行为并完成一次可运行的连接器部署。连接器概述Pulsar 消息如何进入 MongoDBMongoDB sink 是 Pulsar IO 框架中Sinkbyte[]类型的一个实现。从源码注册信息pulsar-io.yaml可以看到该连接器同时提供 source 与 sink 两种形态name: mongo description: MongoDB source and sink connector sinkClass: org.apache.pulsar.io.mongodb.MongoSink sourceClass: org.apache.pulsar.io.mongodb.MongoSource sourceConfigClass: org.apache.pulsar.io.mongodb.MongoConfig sinkConfigClass: org.apache.pulsar.io.mongodb.MongoConfig作为 sink 使用时其数据流为Pulsar topic (消息字节) -- MongoSink.write(Recordbyte[]) -- 内存批量缓存 -- flush() 解析 JSON 为 BSON Document -- collection.insertMany(...) -- MongoDB collection需要特别注意的是该 sink 的类注释明确说明它“假定输入是 JSON 文档”MongoSink.java。每条消息在写入前会被Document.parse()解析为 MongoDB 文档因此发送到 topic 的消息体必须是合法 JSON无法解析的消息会被直接标记为失败record.fail()不会进入 MongoDB。配置项详解MongoDB sink 的配置类为MongoConfigMongoConfig.java官方文档给出了以下属性表其中batchSize与batchTimeMs的默认值在源码中通过常量DEFAULT_BATCH_SIZE 100、DEFAULT_BATCH_TIME_MS 1000定义NameTypeRequiredDefaultDescriptionmongoUriStringtrue空字符串连接器连接 MongoDB 所用的 URI遵循 MongoDB 官方 connection string URI 格式。databaseStringtrue空字符串collection 所属的数据库名称。collectionStringtrue空字符串连接器写入消息的目标 collection 名称。batchSizeintfalse100批量写入 MongoDB 的批大小即攒够多少条消息触发一次写入。batchTimeMslongfalse1000批量操作的触发间隔单位为毫秒。各配置项的行为细节mongoUri必填标准 MongoDB 连接串例如单节点mongodb://localhost:27017、带认证与副本集的多主机形式等。在open()阶段源码通过MongoClients.create(mongoConfig.getMongoUri())创建响应式reactive streamsMongoDB 客户端MongoSink.java因此该 URI 支持的语法与官方 MongoDB 连接串规范一致。底层驱动为mongodb-driver-reactivestreams4.1.2见 pom.xml。database/collection必填sink 打开连接后执行mongoClient.getDatabase(database).getCollection(collection)获取目标集合MongoSink.java。在validate(true, true)中dbRequired与collectionRequired两个参数均为true意味着 sink 场景下二者缺一不可若任一为空会抛出IllegalArgumentException(Required property not set.)。相较之下source 场景对这两个字段的要求有所不同database在源码FieldDoc注释中标注为“source 必须监听sink 必填”。batchSize/batchTimeMs二者共同控制写入节奏详见下文“批量写入机制”。配置校验要求二者必须为正数否则分别抛出batchSize must be a positive integer.与batchTimeMs must be a positive long.MongoConfig.java。配置文件编写JSON 与 YAML在部署 Mongo sink 之前需要通过以下两种方式之一创建配置文件。配置类的加载逻辑位于 MongoConfig.javaload(String yamlFile)使用 Jackson YAML 工厂解析本地 YAML 文件load(MapString, Object)则用于接收框架传入的键值映射如命令行--sink-config参数二者殊途同归。JSON 格式示例{ configs: { mongoUri: mongodb://localhost:27017, database: pulsar, collection: messages, batchSize: 2, batchTimeMs: 500 } }YAML 格式示例官方文档给出的 YAML 示例采用了花括号包裹的写法实际使用时建议使用如下等价的、严格合法的 YAML 映射形式与仓库测试资源 mongoSinkConfig.yaml 保持一致便于直接复制运行mongoUri: mongodb://localhost:27017 database: pulsar collection: messages batchSize: 2 batchTimeMs: 500提示YAML 中的batchSize、batchTimeMs应写为数字类型int/long因为 Jackson 反序列化时会按目标字段类型转换。将上面的配置保存为mongo-sink-config.yaml后即可在下面的部署命令中通过--sink-config-file引用。批量写入机制与消息确认语义理解 Mongo sink 的批处理行为有助于合理设置batchSize与batchTimeMs。其核心逻辑在 MongoSink.java 中由两条触发路径组成路径一按数量触发。每条消息到达后write(Recordbyte[])将记录加入内存列表incomingList当列表大小恰好等于batchSize时立即向 flush 线程池提交一次flush()MongoSink.java。路径二按时间触发。在open()时创建单线程调度器flushExecutor并以batchTimeMs为周期执行scheduleAtFixedRate(() - flush(), batchTimeMs, batchTimeMs, MILLISECONDS)MongoSink.java确保低流量场景下消息也不会在内存中滞留超过一个周期。flush 阶段的处理顺序MongoSink.java加锁取出当前批次的全部记录并重置incomingList逐条将消息字节以 UTF-8 解码再通过Document.parse()解析为 BSONDocument解析失败JsonParseException或BSONException的记录立即调用record.fail()并从批次中剔除对剩余文档一次性执行collection.insertMany(docsToInsert)。确认/失败语义由内部订阅者DocsToInsertSubscriber决定MongoSink.java写入全部成功对批次内所有记录调用record.ack()抛出MongoBulkWriteException批量写部分失败利用getWriteErrors()返回的失败索引仅对未写入成功的记录调用fail()成功部分照常ack()实现部分成功场景下的精确确认其他异常整批全部标记失败。这一行为在单元测试 MongoSinkTest.java 中得到了验证testWriteMultipleMessages模拟 3 条消息中 1 条写入失败最终断言ack()被调用 2 次、fail()被调用 1 次。此外testWriteBadMessage验证了非 JSON 消息如Oops会被直接 fail 而不会写入 MongoDB。部署与运行MongoDB sink 随 Pulsar 发行版以 NAR 归档形式提供模块pulsar-io-mongo版本与仓库一致见 pom.xml。该 NAR 被打包在pulsar-io/目录下文件名为pulsar-io-mongo-version.nar可通过mvn -pl pulsar-io/mongo -am package构建或直接使用官方二进制发行包中预构建的 NAR。本地运行localrun在开发调试阶段可使用pulsar-admin sinks localrun在本地进程内运行 sink。CLI 参数定义于 CmdSinks.java常用参数如下--tenant/--namespace/--namesink 的租户、命名空间与名称--inputs要消费的 Pulsar topic 列表逗号分隔--archiveNAR 归档路径支持本地路径file://或http(s)://URL--sink-config-file指定 YAML 格式的 sink 配置文件路径--sink-config以keyvalue,key2value2形式直接内联传入配置二选一--parallelismsink 实例数--processing-guarantees投递语义at-least-once / at-most-once / effectively-once。示例命令$ pulsar-admin sinks localrun \ --tenant public \ --namespace default \ --name mongo-sink \ --inputs persistent://public/default/mongo-input \ --archive pulsar-io-mongo-2.10.6-SNAPSHOT.nar \ --sink-config-file mongo-sink-config.yaml前提本地需有一个可访问的 MongoDB 实例mongoUri指向的地址以及可消费的 Pulsar 集群或 standalone 实例。集群模式创建确认配置与 NAR 无误后可通过pulsar-admin sinks create将 sink 提交到 Functions Worker 集群运行$ pulsar-admin sinks create \ --tenant public \ --namespace default \ --name mongo-sink \ --inputs persistent://public/default/mongo-input \ --archive pulsar-io-mongo-2.10.6-SNAPSHOT.nar \ --sink-config-file mongo-sink-config.yaml \ --parallelism 1创建后可用pulsar-admin sinks status --tenant public --namespace default --name mongo-sink查看运行状态。关于 sink 的更完整 CLI 选项可进一步查阅 io-cli.md 与 io-use.md。测试与验证仓库为 MongoDB 连接器提供了较完整的单元测试可作为理解行为边界的参考MongoConfigTest.java验证load(map)与load(yamlFile)两种加载方式的字段映射同时用expectedExceptionsMessageRegExp断言四种非法配置缺少必填项、batchSize/batchTimeMs非正会抛出预期的IllegalArgumentExceptionMongoSinkTest.java通过 Mockito 注入 mock 的MongoClient覆盖正常写入 ack、空消息 fail、批量部分失败 ack/fail 混合、非 JSON 消息 fail 等场景mongoSinkConfig.yaml测试用配置样例与测试常量TestHelper.java中的URImongodb://localhost、DBpulsar、COLLmessages、BATCH_SIZE2、BATCH_TIME500完全对应。小结MongoDB sink connector 是典型的“配置即服务”式连接器只需提供mongoUri、database、collection三个必填项即可把 Pulsar topic 中的 JSON 消息持续、批量地落入 MongoDB。其内部通过“数量阈值 时间周期”双触发机制控制写入频率并利用insertMany与精确的 ack/fail 索引映射在批量写入部分失败时仍能保证 Pulsar 消息确认语义的正确性。合理配置batchSize与batchTimeMs默认分别为 100 与 1000ms可以在吞吐与延迟之间取得平衡——例如高吞吐场景可调大batchSize对延迟敏感的场景则可调小batchTimeMs。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar HDFS3 Sink Connector 完全指南从配置到源码级原理Apache Pulsar HDFS3 Sink Connector 完全指南从配置到源码级原理 HDFS3 sink connector 是 Apache消息队列后端流处理Apache Pulsar Kafka Sink Connector 完全指南配置、运行与源码原理Apache Pulsar Kafka Sink Connector 完全指南配置、运行与源码原理 Kafka sink connector 是 Apache消息队列后端流处理Apache Pulsar RabbitMQ Sink Connector 完全指南配置、部署与源码级原理剖析Apache Pulsar RabbitMQ Sink Connector 完全指南配置、部署与源码级原理剖析 Apache Pulsar 的 RabbitM消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

Vivado联合Questasim仿真流程
Vivado联合Questasim仿真流程

1、Settings设置Questasim仿真:2、Export Simutalion点击后出现下图所设设置,配置Export directory目录:设置完毕后点击OK,会在Export Directory设置的目录中出现该目录:queats目录中内容如下:导出的仿真脚… · 2026/9/24 20:25:26

OpenSandbox实战:轻量级进程隔离沙箱的部署与配置指南
OpenSandbox实战:轻量级进程隔离沙箱的部署与配置指南

我最近在几个开发环境里反复折腾应用隔离的方案,最后被一个叫 OpenSandbox 的命令行工具给留住了。这东西说白了就是一个开源的应用级沙箱运行环境,能把不太可信的脚本、二进制程序、甚至整组服务进程关进一个受限的运行空间里,让它在里面折腾… · 2026/9/24 20:25:20

OpenSandbox极简部署与实践:让不可信代码在隔离沙盒中安全运行
OpenSandbox极简部署与实践:让不可信代码在隔离沙盒中安全运行

1. OpenSandbox到底解决什么问题:从一次重装系统的教训说起1.1 一个让人崩溃的开发场景先说我自己的经历。去年有段时间,我在研究一个第三方提供的自动化测试脚本,对方打包了一堆二进制文件和一个安装入口,文档里写着“建议在干净… · 2026/9/24 20:25:20

wardogs联机卡顿掉线怎么办?从网络排查到优化的完整指南
wardogs联机卡顿掉线怎么办?从网络排查到优化的完整指南

大家联机打wardogs遇到卡顿掉线,不管是组队冲锋还是和认识的朋友开黑,体验确实会大打折扣。我玩这游戏也算有些年头了,早期在宿舍用公共WiFi,高峰期几乎走两步就漂移,后来换到有线网络加上一系列排查调整,才… · 2026/9/24 21:33:03

代号鸢联机掉线不用慌,从手机到光猫的排查指南
代号鸢联机掉线不用慌,从手机到光猫的排查指南

玩了这么久代号鸢,相信不少人都有过这种体验:正打到关键回合,或者剧情推到高潮,屏幕一卡,然后眼睁睁看着角色原地发呆,接着就是“连接已断开”的提示。尤其是在联机打地宫、打首领的环节,卡顿和… · 2026/9/24 21:33:03

DeepSeek Harness评测:多智能体协作编排与配置实战
DeepSeek Harness评测:多智能体协作编排与配置实战

DeepSeek官方这次的低调操作确实让我有点意外,要不是在官网角落看到更新日志,我根本不知道出了一个叫Harness的桌面端。用了一周多,体感很新鲜,和之前那些套壳对话客户端完全不是一回事。这篇就把我的实测记录、踩坑经历、还有一些… · 2026/9/24 21:33:03

任务向量合并:从参数平均到映射优化
任务向量合并:从参数平均到映射优化

1. 任务向量这东西,加法本来就不该是默认操作我第一次在论文里看到任务向量(Task Vector)时,说实话是带着怀疑的。把模型在目标任务上做微调,然后拿微调后的权重减去预训练权重,得到一个“向量”&#xff0… · 2026/9/24 21:33:03

HarmonyOS证书过期怎么办?从紧急修复到系统化管理全攻略
HarmonyOS证书过期怎么办?从紧急修复到系统化管理全攻略

HarmonyOS开发做到一定阶段,证书过期这件事大概率会找上门。它不像编译报错那样当场给你红字提示,更多时候是悄无声息地卡住你的调试、阻断你的上架,甚至让线上推送直接哑火。我在多个项目里处理过HarmonyOS应用从调试到发布的完整流程&#… · 2026/9/24 21:33:03

FastJson核心要点与避坑:序列化、时间处理与安全升级
FastJson核心要点与避坑:序列化、时间处理与安全升级

1. 为什么 FastJson 至今仍是Java后端绕不开的JSON库讲真,我最早接触 FastJson 是在 2016 年,那时候公司的老项目里到处都是JSON.toJSONString()和JSON.parseObject(),后来自己写新服务,也跟着用了几年。中途经历了几次 FastJson … · 2026/9/24 21:32:57

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为… · 2026/9/24 0:00:13

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码