ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

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),仅供参考
返回列表