news 2026/9/23 9:32:26

Apache Pulsar Connector 调试实战指南:localrun 与集群模式下的日志、Admin CLI 与排错清单

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar Connector 调试实战指南:localrun 与集群模式下的日志、Admin CLI 与排错清单

Apache Pulsar Connector 调试实战指南:localrun 与集群模式下的日志、Admin CLI 与排错清单

【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar

本指南以 Apache Pulsar 的 Mongo sink 连接器为贯穿示例,系统讲解在 localrun 与 cluster(集群)两种模式下调试 Source/Sink 连接器的完整方法,包括调试环境搭建、连接器日志的逐段解读、pulsar-admin管理命令(get/status/topics stats)的使用,以及一份可复用的连接器调试检查清单。读完本文,你将能够独立定位连接器运行失败、配置错误、消息未写入外部系统等常见问题。

调试之前:理解连接器的运行形态

Pulsar 连接器(connector)分为 Source(从外部系统读取数据写入 Pulsar)和 Sink(从 Pulsar 读取消息写入外部系统)两类,它们在底层都作为 Function 运行在 Functions worker 上。这意味着连接器的生命周期、运行时、日志与函数(Function)完全一致——这一点可以从 site2/docs/io-overview.md 中"Connectors (sources and sinks) and Functions are components of instances, and they all run on Functions workers"的描述得到印证。

连接器有两种启动方式,调试手段也因此分为两条主线:

  • localrun 模式:连接器作为本地进程/线程在运行pulsar-admin命令的机器上启动,不经过 Functions worker 调度,日志直接输出到控制台,适合快速验证和单机排查。
  • cluster 模式:通过pulsar-admin sinks create/sources create将连接器提交到 Functions worker 集群,由 worker 分配实例运行,日志落盘到 worker 所在节点,需配合 Admin CLI 远程查询状态。

准备调试环境:以 Mongo sink 为例

为了演示完整的调试过程,需要先部署一套可复现的最小环境:一个 Mongo 服务 + 一个 Pulsar standalone 实例 + Mongo sink 的 nar 包。

1. 启动 Mongo 服务

使用 Docker 拉取并启动 Mongo 4,将数据目录挂载到宿主机$PWD/data,映射 27017 端口:

docker pull mongo:4 docker run -d -p 27017:27017 --name pulsar-mongo -v $PWD/data:/data/db mongo:4

2. 创建数据库与集合

进入容器,通过mongoshell 创建名为pulsar的数据库和名为messages的集合:

docker exec -it pulsar-mongo /bin/bash mongo > use pulsar > db.createCollection('messages') > exit

3. 启动 Pulsar standalone

拉取apachepulsar/pulsar:2.4.0镜像,以--link pulsar-mongo与 Mongo 容器连通,映射 6650(Broker 端口)和 8080(Web 服务端口):

docker pull apachepulsar/pulsar:2.4.0 docker run -d -it -p 6650:6650 -p 8080:8080 -v $PWD/data:/pulsar/data --link pulsar-mongo --name pulsar-mongo-standalone apachepulsar/pulsar:2.4.0 bin/pulsar standalone

4. 准备 Mongo sink 配置文件

新建mongo-sink-config.yaml,指定 Mongo 连接串、目标库表与批量写入参数:

configs: mongoUri: "mongodb://pulsar-mongo:27017" database: "pulsar" collection: "messages" batchSize: 2 batchTimeMs: 500

将配置文件拷入 Pulsar 容器:

docker cp mongo-sink-config.yaml pulsar-mongo-standalone:/pulsar/

参数说明(对应源码 pulsar-io/mongo/src/main/java/org/apache/pulsar/io/mongodb/MongoConfig.java):

  • mongoUri:MongoDB 连接串,必填项,格式参考 MongoDB 官方 connection string 文档;
  • database:目标数据库名,Sink 模式下必填;
  • collection:消息写入的目标集合名,Sink 模式下必填;
  • batchSize:批量写入的条数阈值,默认值为 100(DEFAULT_BATCH_SIZE);
  • batchTimeMs:批量写入的时间间隔(毫秒),默认值为 1000(DEFAULT_BATCH_TIME_MS)。

从源码中的validate()方法可以看到,mongoUridatabasecollection任一为空都会抛出IllegalArgumentException("Required property not set."),而batchSize必须为正整数、batchTimeMs必须为正 long,否则启动即失败——这正是调试时要重点核对的前置条件。

5. 下载 Mongo sink 的 nar 包

进入 Pulsar 容器,下载对应版本的连接器 nar 包:

docker exec -it pulsar-mongo-standalone /bin/bash curl -O http://apache.01link.hk/pulsar/pulsar-2.4.0/connectors/pulsar-io-mongo-2.4.0.nar

提示:nar 包是 Pulsar 连接器/函数的打包格式,内部包含连接器类以及META-INF/bundled-dependencies/下的全部依赖。当前仓库中 Mongo sink 的源码位于 pulsar-io/mongo/src/main/java/org/apache/pulsar/io/mongodb/MongoSink.java,其@Connector注解声明了连接器名mongo、类型SINK与配置类MongoConfig,这些元信息在 nar 包加载与命令行校验时都会被用到。

在 localrun 模式下调试

启动 localrun

localrun 模式使用pulsar-admin sinks localrun命令启动连接器,它把连接器直接运行在当前节点上(而不是提交给 worker 集群),非常适合调试。关于localrun命令的完整参数,可参考 site2/docs/reference-connector-admin.md。

./bin/pulsar-admin sinks localrun \ --archive connectors/pulsar-io-mongo-{{pulsar:version}}.nar \ --tenant public --namespace default \ --inputs test-mongo \ --name pulsar-mongo-sink \ --sink-config-file mongo-sink-config.yaml \ --parallelism 1

从源码看,localrun的实现类是 pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java。它内部会根据运行环境自动选择两种运行时:

  • ThreadRuntime(线程模式):默认方式,连接器作为 Java 线程在本地 JVM 中运行(对应源码中的startThreadedModeThreadRuntimeFactory);
  • ProcessRuntime(进程模式):连接器作为独立子进程运行(对应startProcessModeProcessRuntimeFactory)。

localrun启动后,连接器从指定 topic(test-mongo)消费消息,由MongoSink实例将消息解析为 BSON 文档并批量写入 Mongo。

使用连接器日志

localrun 模式下获取连接器日志有两种途径:

  1. 控制台日志:执行localrun命令后,日志会自动打印到终端控制台;

  2. 日志文件:日志同时写入文件,路径规律为:

    logs/functions/tenant/namespace/function-name/function-name-instance-id.log

    以本例的 Mongo sink 为例,日志文件位于:

    logs/functions/public/default/pulsar-mongo-sink/pulsar-mongo-sink-0.log

    说明:logs是 Pulsar standalone 运行目录下的日志根目录;路径中的tenant/namespace对应--tenant public --namespace defaultfunction-name对应--name pulsar-mongo-sinkinstance-id对应实例编号(并行度为 1 时只有实例 0)。

日志逐段解读

连接器启动日志信息量大,下面将其拆分为小块并逐一解释其含义与调试价值。

第一段:nar 包解压路径

08:21:54.132 [main] INFO org.apache.pulsar.common.nar.NarClassLoader - Created class loader with paths: [file:/tmp/pulsar-nar/pulsar-io-mongo-2.4.0.nar-unpacked/, file:/tmp/pulsar-nar/pulsar-io-mongo-2.4.0.nar-unpacked/META-INF/bundled-dependencies/,

这段日志说明NarClassLoader已经成功创建类加载器,并给出了 nar 包解压后的存储路径。LocalRunner中通过narExtractionDirectory指定 nar 的解压目录(默认是/tmp/pulsar-nar下的临时目录),解压后的META-INF/bundled-dependencies/存放连接器的全部第三方依赖。

调试提示:如果抛出了class cannot be found(类找不到)异常,请检查 nar 包是否被完整解压到file:/tmp/pulsar-nar/pulsar-io-mongo-2.4.0.nar-unpacked/META-INF/bundled-dependencies/目录中。该异常通常意味着 nar 包损坏、下载不完整,或打包时依赖缺失。

第二段:连接器实例配置(InstanceConfig)

08:21:55.390 [main] INFO org.apache.pulsar.functions.runtime.ThreadRuntime - ThreadContainer starting function with instance config InstanceConfig(instanceId=0, functionId=853d60a1-0c48-44d5-9a5c-6917386476b2, functionVersion=c2ce1458-b69e-4175-88c0-a0a856a2be8c, functionDetails=tenant: "public" namespace: "default" name: "pulsar-mongo-sink" className: "org.apache.pulsar.functions.api.utils.IdentityFunction" autoAck: true parallelism: 1 source { typeClassName: "[B" inputSpecs { key: "test-mongo" value { } } cleanupSubscription: true } sink { className: "org.apache.pulsar.io.mongodb.MongoSink" configs: "{\"mongoUri\":\"mongodb://pulsar-mongo:27017\",\"database\":\"pulsar\",\"collection\":\"messages\",\"batchSize\":2,\"batchTimeMs\":500}" typeClassName: "[B" } resources { cpu: 1.0 ram: 1073741824 disk: 10737418240 } componentType: SINK , maxBufferedTuples=1024, functionAuthenticationSpec=null, port=38459, clusterName=local)

这条日志由ThreadRuntime打印(对应源码 pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/thread/ThreadRuntime.java 中的log.info("ThreadContainer starting function with instanceId {} functionId {} namespace {}"...))。它完整展示了连接器的实例配置,是核对连接器是否配置正确的第一手依据:

  • tenant/namespace/name:连接器的三元组标识,应与你启动时传入的参数一致;
  • className(source 侧):此处为IdentityFunction,即 Sink 的输入源是一个透传函数,不做数据变换;
  • sink.classNameorg.apache.pulsar.io.mongodb.MongoSink,即实际执行的 Sink 实现类;
  • sink.configs:JSON 序列化后的mongo-sink-config.yaml内容,mongoUridatabasecollectionbatchSize=2batchTimeMs=500全部可见;
  • parallelism: 1resources:并行度与 CPU/内存/磁盘资源配额;
  • maxBufferedTuples=1024:每个实例允许缓冲的最大消息数;
  • clusterName=local:localrun 模式下的固定集群标识。

若这里的configs内容与你预期不符(例如database拼写错误、batchSize 为 0 或负数),基本可以断定是配置文件问题,应回到mongo-sink-config.yaml检查。因为MongoConfig.validate()会在配置非法时直接抛异常,连接器将无法完成启动。

第三段:Mongo 连接状态

08:21:56.231 [cluster-ClusterId{value='5d6396a3c9e77c0569ff00eb', description='null'}-pulsar-mongo:27017] INFO org.mongodb.driver.connection - Opened connection [connectionId{localValue:1, serverValue:8}] to pulsar-mongo:27017 08:21:56.326 [cluster-ClusterId{value='5d6396a3c9e77c0569ff00eb', description='null'}-pulsar-mongo:27017] INFO org.mongodb.driver.cluster - Monitor thread successfully connected to server with description ServerDescription{address=pulsar-mongo:27017, type=STANDALONE, state=CONNECTED, ok=true, version=ServerVersion{versionList=[4, 2, 0]}, minWireVersion=0, maxWireVersion=8, maxDocumentSize=16777216, logicalSessionTimeoutMinutes=30, roundTripTimeNanos=89058800}

这两条日志来自 Mongo Java Driver:第一条说明与pulsar-mongo:27017建立了 TCP 连接,第二条说明集群监控线程成功连接服务器,ServerDescription中给出了服务器类型(STANDALONE)、状态(CONNECTED)、ok=true、Mongo 版本(4.2.0)等关键信息。

调试价值:如果这里出现连接失败或超时,说明mongoUri指向的地址不可达(例如容器名解析失败、端口未映射、Mongo 未启动)。注意此处主机名是pulsar-mongo(Docker 容器名),若脱离容器环境运行,需要改为实际可达的 IP 或域名。

第四段:Consumer 与 Client 配置

08:21:56.719 [pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerStatsRecorderImpl - Starting Pulsar consumer status recorder with config: { "topicNames" : [ "test-mongo" ], "topicsPattern" : null, "subscriptionName" : "public/default/pulsar-mongo-sink", "subscriptionType" : "Shared", "receiverQueueSize" : 1000, "acknowledgementsGroupTimeMicros" : 100000, "negativeAckRedeliveryDelayMicros" : 60000000, "maxTotalReceiverQueueSizeAcrossPartitions" : 50000, "consumerName" : null, "ackTimeoutMillis" : 0, "tickDurationMillis" : 1000, "priorityLevel" : 0, "cryptoFailureAction" : "CONSUME", "properties" : { "application" : "pulsar-sink", "id" : "public/default/pulsar-mongo-sink", "instance_id" : "0" }, "readCompacted" : false, "subscriptionInitialPosition" : "Latest", "patternAutoDiscoveryPeriod" : 1, "regexSubscriptionMode" : "PersistentOnly", "deadLetterPolicy" : null, "autoUpdatePartitions" : true, "replicateSubscriptionState" : false, "resetIncludeHead" : false } 08:21:56.726 [pulsar-client-io-1-1] INFO org.apache.pulsar.client.impl.ConsumerStatsRecorderImpl - Pulsar client config: { "serviceUrl" : "pulsar://localhost:6650", "authPluginClassName" : null, "authParams" : null, "operationTimeoutMs" : 30000, "statsIntervalSeconds" : 60, "numIoThreads" : 1, "numListenerThreads" : 1, "connectionsPerBroker" : 1, "useTcpNoDelay" : true, "useTls" : false, "tlsTrustCertsFilePath" : null, "tlsAllowInsecureConnection" : false, "tlsHostnameVerificationEnable" : false, "concurrentLookupRequest" : 5000, "maxLookupRequest" : 50000, "maxNumberOfRejectedRequestPerConnection" : 50, "keepAliveIntervalSeconds" : 30, "connectionTimeoutMs" : 10000, "requestTimeoutMs" : 60000, "defaultBackoffIntervalNanos" : 100000000, "maxBackoffIntervalNanos" : 30000000000 }

这两段日志分别输出Consumer 配置Pulsar Client 配置,用于核对消费侧行为:

  • Consumer 配置:消费的 topic 为test-mongo;订阅名为public/default/pulsar-mongo-sink(订阅名 = 三元组标识,固定规则);订阅类型为Shared(共享订阅,多个实例可并行消费);receiverQueueSize=1000为接收队列大小;subscriptionInitialPosition=Latest表示新订阅默认从最新消息开始消费;acknowledgementsGroupTimeMicros=100000表示 100ms 的确认聚合窗口;negativeAckRedeliveryDelayMicros=60000000表示负确认后 60s 重投递。属性中application=pulsar-sinkid=public/default/pulsar-mongo-sinkinstance_id=0会同步体现在 topic stats 的消费者元数据里。
  • Client 配置serviceUrl=pulsar://localhost:6650是 localrun 模式连接 Broker 的服务地址;numIoThreads=1numListenerThreads=1为 IO/监听线程数;connectionsPerBroker=1operationTimeoutMs=30000keepAliveIntervalSeconds=30;未开启 TLS(useTls=false)。

调试价值:若连接器"启动成功但收不到消息",优先核对此处serviceUrl是否正确、subscriptionInitialPosition是否为预期位置(Latest会跳过此前已发布的历史消息)。需要消费历史消息时,可通过--subs-position类参数调整订阅初始位置。

在集群(cluster)模式下调试

cluster 模式下连接器被提交到 Functions worker 运行,调试手段分为两类:连接器日志Admin CLI

使用连接器日志

集群模式下,一个 worker 上可能同时运行多个连接器实例。要定位某个连接器的日志路径,需要使用workerId(worker 实例的唯一标识)来确定该连接器运行在哪个 worker 节点上,再前往该节点的logs/functions/tenant/namespace/function-name/目录查看对应实例的日志文件。workerId的获取方式见下文status命令。

使用 Admin CLI

Pulsar Admin CLI 提供了三个与调试直接相关的子命令:getstatustopics stats

首先,在集群模式下创建 Mongo sink:

./bin/pulsar-admin sinks create \ --archive pulsar-io-mongo-2.4.0.nar \ --tenant public \ --namespace default \ --inputs test-mongo \ --name pulsar-mongo-sink \ --sink-config-file mongo-sink-config.yaml \ --parallelism 1

get:获取连接器基础信息

get命令返回连接器的完整配置快照(tenant、namespace、name、parallelism、className、configs 等),用于确认连接器是否按预期配置创建get命令的更多选项见 site2/docs/reference-connector-admin.md。

./bin/pulsar-admin sinks get --tenant public --namespace default --name pulsar-mongo-sink { "tenant": "public", "namespace": "default", "name": "pulsar-mongo-sink", "className": "org.apache.pulsar.io.mongodb.MongoSink", "inputSpecs": { "test-mongo": { "isRegexPattern": false } }, "configs": { "mongoUri": "mongodb://pulsar-mongo:27017", "database": "pulsar", "collection": "messages", "batchSize": 2.0, "batchTimeMs": 500.0 }, "parallelism": 1, "processingGuarantees": "ATLEAST_ONCE", "retainOrdering": false, "autoAck": true }

输出要点解读:

  • className:实际运行的 Sink 实现类,应为org.apache.pulsar.io.mongodb.MongoSink
  • configs:配置文件中所有参数(batchSizebatchTimeMs以数值形式返回);
  • parallelism:实例并行度,此处为 1;
  • processingGuarantees:处理保证级别,ATLEAST_ONCE表示至少一次语义(MongoSink内部采用批量写入 + 成功/失败分别ack/fail的机制,与这一语义一致);
  • retainOrdering:是否保持消息顺序,false表示不做顺序保证;
  • autoAck:是否自动确认,true表示由框架自动完成确认。

get返回的configs与配置文件不一致,说明创建命令传入的参数有问题,应重新核对--sink-config-file指向的文件内容。

status:获取连接器运行状态

status命令返回实例数量、运行中实例数、每个实例的instanceIdworkerId以及各类错误/统计计数,是判断连接器是否健康运行的核心命令。更多选项见 site2/docs/reference-connector-admin.md。

./bin/pulsar-admin sinks status --tenant public \ --namespace default \ --name pulsar-mongo-sink { "numInstances" : 1, "numRunning" : 1, "instances" : [ { "instanceId" : 0, "status" : { "running" : true, "error" : "", "numRestarts" : 0, "numReadFromPulsar" : 0, "numSystemExceptions" : 0, "latestSystemExceptions" : [ ], "numSinkExceptions" : 0, "latestSinkExceptions" : [ ], "numWrittenToSink" : 0, "lastReceivedTime" : 0, "workerId" : "c-standalone-fw-5d202832fd18-8080" } } ] }

输出要点解读:

  • numInstances/numRunning:期望实例数与实际运行实例数。若numRunning < numInstances,说明有实例启动失败或反复重启;
  • running:当前实例是否在运行;
  • error:实例级错误信息,为空表示无错误;
  • numRestarts:实例重启次数,数值持续增长说明连接器不稳定(如配置非法反复失败、OOM、外部依赖不可用);
  • numSystemExceptions/latestSystemExceptions:系统级异常计数与最近异常明细(如 Pulsar client 连接异常、配置解析异常);
  • numSinkExceptions/latestSinkExceptions:Sink 写外部系统时的异常计数与最近异常明细(如 Mongo 写入失败);
  • numReadFromPulsar/numWrittenToSink:从 Pulsar 读取的消息数与成功写入 Sink 的消息数,两者长期为 0 或严重不对称时都值得深挖;
  • lastReceivedTime:最近一次收到消息的时间戳,为 0 表示从未收到过消息;
  • workerId实例所在 worker 的唯一标识。若多个连接器运行在同一 worker 上,workerId可帮你定位该连接器运行在哪个节点,从而找到对应的日志文件。

调试价值:latestSystemExceptionslatestSinkExceptions是最直接的故障线索来源。Sink 写外部系统失败(例如 Mongo 认证失败、目标集合不存在)时,numSinkExceptions会增长且latestSinkExceptions会携带堆栈信息。

topics stats:获取主题与消费者统计

topics stats命令返回指定 topic 及其生产者和消费者的统计信息,用于判断**消息是否到达 topic、是否存在积压(backlog)、消费者是否有可用许可(permits)**等。所有速率指标基于 1 分钟窗口计算,相对上一个完整 1 分钟周期而言。

./bin/pulsar-admin topics stats test-mongo { "msgRateIn" : 0.0, "msgThroughputIn" : 0.0, "msgRateOut" : 0.0, "msgThroughputOut" : 0.0, "averageMsgSize" : 0.0, "storageSize" : 1, "publishers" : [ ], "subscriptions" : { "public/default/pulsar-mongo-sink" : { "msgRateOut" : 0.0, "msgThroughputOut" : 0.0, "msgRateRedeliver" : 0.0, "msgBacklog" : 0, "blockedSubscriptionOnUnackedMsgs" : false, "msgDelayed" : 0, "unackedMessages" : 0, "type" : "Shared", "msgRateExpired" : 0.0, "consumers" : [ { "msgRateOut" : 0.0, "msgThroughputOut" : 0.0, "msgRateRedeliver" : 0.0, "consumerName" : "dffdd", "availablePermits" : 999, "unackedMessages" : 0, "blockedConsumerOnUnackedMsgs" : false, "metadata" : { "instance_id" : "0", "application" : "pulsar-sink", "id" : "public/default/pulsar-mongo-sink" }, "connectedSince" : "2019-08-26T08:48:07.582Z", "clientVersion" : "2.4.0", "address" : "/172.17.0.3:57790" } ], "isReplicated" : false } }, "replication" : { }, "deduplicationStatus" : "Disabled" }

输出要点解读:

  • topic 整体msgRateIn/msgThroughputIn(生产速率)、msgRateOut/msgThroughputOut(消费速率)、storageSize(存储大小)、publishers(生产者列表)。若publishers为空且msgRateIn=0,说明还没有生产者向该 topic 发布消息;
  • 订阅public/default/pulsar-mongo-sinkmsgBacklog(积压消息数)、unackedMessages(未确认消息数)、msgRateRedeliver(重投递速率)、type=Shared(订阅类型)。msgBacklog持续增长而msgRateOut=0,说明消费者没有有效消费,应回到status检查 Sink 异常;unackedMessages长期不降且blockedSubscriptionOnUnackedMsgs=true,说明存在确认堆积问题;
  • 消费者明细availablePermits(可用许可数,接近receiverQueueSize=1000表示消费者空闲)、metadata中的application=pulsar-sinkid=public/default/pulsar-mongo-sink可确认该消费者正是本连接器的实例、connectedSince(连接建立时间)、clientVersionaddress(消费者地址,可用于与status中的workerId对应)。

调试价值:topics stats能够把"消息流"这一环节单独隔离出来验证——如果 topic 侧一切正常(有生产、有消费、无积压、无异常),但外部系统里没有数据,那么问题很可能出在 Sink 实现或外部系统一侧,需要结合连接器日志与外部系统日志继续排查。

连接器调试检查清单(Checklist)

以下清单覆盖连接器调试时需要核对的全部关键区域,既可作为全面排查的提醒,也可作为评估连接器当前状态的工具:

  1. Pulsar 是否正常启动?检查 Broker、BookKeeper、ZooKeeper 以及 standalone 模式下的 Web 服务是否可用,bin/pulsar-admin brokers healthcheck类命令可用于快速验证。
  2. 外部服务是否正常运行?确认目标系统(本例为 Mongo)进程存活、端口可连通、凭据与权限正确,必要时在外部系统侧查看其自身日志。
  3. nar 包是否完整?检查 nar 文件是否下载完整、版本与 Pulsar 版本匹配,解压目录META-INF/bundled-dependencies/中依赖是否齐全(缺依赖通常表现为class cannot be found)。
  4. 连接器配置文件是否正确?对照MongoConfig的必填项与取值范围逐项核对mongoUridatabasecollectionbatchSizebatchTimeMs,并结合get命令返回的configs做双重确认。
  5. localrun 模式:运行连接器并检查控制台打印的连接器日志(nar 解压路径、InstanceConfig、Mongo 连接、Consumer/Client 配置四段信息逐一核对)。
  6. cluster 模式
    • get命令获取基础配置信息;
    • status命令获取实例运行状态、异常计数与workerId
    • topics stats命令获取 topic 及其生产/消费者的统计信息;
    • 根据workerId定位 worker 节点,检查连接器日志文件。
  7. 进入外部系统验证结果:登录 Mongo,检查pulsarmessages集合中是否出现符合预期的文档,从而确认端到端链路(Pulsar → Sink → Mongo)是否真正打通。

从源码理解调试的关键机制

localrun 的两类运行时

LocalRunner(pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java)根据配置选择线程模式或进程模式:

  • 线程模式(startThreadedModeThreadRuntimeFactory):连接器实例运行在pulsar-admin进程内部的独立线程组中,日志直接复用当前进程的日志框架,因此会实时打印在控制台,同时也由框架写入日志文件;
  • 进程模式(startProcessModeProcessRuntimeFactory):连接器实例以独立子进程运行,调试时可通过进程输出重定向或日志文件观察其行为。

无论哪种模式,每个实例的InstanceConfig都会记录functionIdfunctionVersioninstanceIdclusterName=localmaxBufferedTuples=1024等信息,这些正是启动日志第二段所展示的内容,也是ThreadRuntime打印ThreadContainer starting function with instanceId {} functionId {} namespace {}日志的数据来源。

MongoSink 的写入与确认语义

MongoSink.java 展示了 Sink 侧的典型实现模式,理解它有助于解读status中的计数指标:

  • open():加载并校验MongoConfig,创建 Mongo 客户端,获取目标库表,并启动一个定时调度线程(flushExecutor),按batchTimeMs周期触发flush()
  • write():每收到一条消息,先将Record加入内存缓冲incomingList;当缓冲条数达到batchSize时,立即触发一次异步flush()
  • flush():将缓冲中的消息逐个解析为 BSONDocument,解析失败(JSON 格式错误)的消息调用record.fail()并剔除;解析成功的批量调用collection.insertMany()
  • 写入结果通过DocsToInsertSubscriber回调处理:全部成功则逐条ack();发生MongoBulkWriteException时根据写入错误索引区分成功与失败的消息,成功者ack()、失败者fail(),从而保证"至少一次"的处理语义。

调试启发:batchSize=2batchTimeMs=500的配置意味着每凑齐 2 条消息或每 500ms 就会批量写入一次。如果statusnumWrittenToSink小于numReadFromPulsar,可优先检查latestSinkExceptions是否出现 BSON 解析错误(消息不是合法 JSON)或 Mongo 写入异常。对应的测试用例位于 pulsar-io/mongo/src/test/java/org/apache/pulsar/io/mongodb/MongoSinkTest.java,可作为理解其行为与复现问题的参考。

小结

调试 Pulsar 连接器的核心方法论可以概括为"由近及远、逐层隔离":先用localrun在本地快速复现并借助控制台日志定位配置与依赖问题;进入集群环境后,用get核对配置、用status观察实例健康与异常明细、用topics stats隔离消息链路问题,再结合workerId定位节点查看日志;最后进入外部系统验证端到端结果。配合本文的调试检查清单,即可系统化地完成连接器的排查与验证。

【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 9:32:00

解决PowerShell启动自动跳转桌面的问题

1. 问题现象与背景解析最近在Windows环境下使用PowerShell Core&#xff08;简称pwsh&#xff09;时&#xff0c;发现一个让人困扰的现象&#xff1a;无论是通过CMD命令行直接启动pwsh&#xff0c;还是在VS Code中新建终端窗口&#xff0c;系统总是会自动跳转到桌面目录。作为一…

作者头像 李华
网站建设 2026/9/23 9:27:11

Python+Django构建高效餐饮管理系统实战

1. 项目概述&#xff1a;餐饮管理系统的数字化转型在餐饮行业竞争日益激烈的今天&#xff0c;一套高效的个性化管理系统已成为门店运营的刚需。我最近用PythonDjango完整开发了一套餐饮管理系统&#xff0c;从点餐、库存到会员管理全覆盖。这个系统特别适合中小型餐饮企业&…

作者头像 李华
网站建设 2026/9/23 9:26:14

【二分查找】LC 33.搜索旋转排序数组

文章目录前言一、题目1、原题链接2、题目描述二、个人思路整理1、思路分析2、解题代码三、知识风暴前言 本专栏文章为《LeetCode 热题 100》的刷题题解&#xff0c;相关内容如有侵权&#xff0c;立即删除。 一、题目 1、原题链接 33.搜索旋转排序数组 2、题目描述 二、个人思路…

作者头像 李华
网站建设 2026/9/23 9:26:13

衍射光束扩散器设计:从原理到工程实践

1. 衍射光束扩散器设计概述在光学工程领域&#xff0c;衍射光束扩散器&#xff08;Diffractive Optical Element, DOE&#xff09;是一种能够将入射激光束转换为特定光场分布的光学元件。最近我在VirtualLab Fusion平台上完成了一个实际项目&#xff1a;设计一个能将公司标识投…

作者头像 李华