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:42. 创建数据库与集合进入容器通过mongoshell 创建名为pulsar的数据库和名为messages的集合docker exec -it pulsar-mongo /bin/bash mongo use pulsar db.createCollection(messages) exit3. 启动 Pulsar standalone拉取apachepulsar/pulsar:2.4.0镜像以--link pulsar-mongo与 Mongo 容器连通映射 6650Broker 端口和 8080Web 服务端口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 standalone4. 准备 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.javamongoUriMongoDB 连接串必填项格式参考 MongoDB 官方 connection string 文档database目标数据库名Sink 模式下必填collection消息写入的目标集合名Sink 模式下必填batchSize批量写入的条数阈值默认值为 100DEFAULT_BATCH_SIZEbatchTimeMs批量写入的时间间隔毫秒默认值为 1000DEFAULT_BATCH_TIME_MS。从源码中的validate()方法可以看到mongoUri、database、collection任一为空都会抛出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 模式下调试启动 localrunlocalrun 模式使用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 中运行对应源码中的startThreadedMode与ThreadRuntimeFactoryProcessRuntime进程模式连接器作为独立子进程运行对应startProcessMode与ProcessRuntimeFactory。localrun启动后连接器从指定 topictest-mongo消费消息由MongoSink实例将消息解析为 BSON 文档并批量写入 Mongo。使用连接器日志localrun 模式下获取连接器日志有两种途径控制台日志执行localrun命令后日志会自动打印到终端控制台日志文件日志同时写入文件路径规律为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 包损坏、下载不完整或打包时依赖缺失。第二段连接器实例配置InstanceConfig08:21:55.390 [main] INFO org.apache.pulsar.functions.runtime.ThreadRuntime - ThreadContainer starting function with instance config InstanceConfig(instanceId0, functionId853d60a1-0c48-44d5-9a5c-6917386476b2, functionVersionc2ce1458-b69e-4175-88c0-a0a856a2be8c, functionDetailstenant: 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 , maxBufferedTuples1024, functionAuthenticationSpecnull, port38459, clusterNamelocal)这条日志由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连接器的三元组标识应与你启动时传入的参数一致classNamesource 侧此处为IdentityFunction即 Sink 的输入源是一个透传函数不做数据变换sink.classNameorg.apache.pulsar.io.mongodb.MongoSink即实际执行的 Sink 实现类sink.configsJSON 序列化后的mongo-sink-config.yaml内容mongoUri、database、collection、batchSize2、batchTimeMs500全部可见parallelism: 1、resources并行度与 CPU/内存/磁盘资源配额maxBufferedTuples1024每个实例允许缓冲的最大消息数clusterNamelocallocalrun 模式下的固定集群标识。若这里的configs内容与你预期不符例如database拼写错误、batchSize 为 0 或负数基本可以断定是配置文件问题应回到mongo-sink-config.yaml检查。因为MongoConfig.validate()会在配置非法时直接抛异常连接器将无法完成启动。第三段Mongo 连接状态08:21:56.231 [cluster-ClusterId{value5d6396a3c9e77c0569ff00eb, descriptionnull}-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{value5d6396a3c9e77c0569ff00eb, descriptionnull}-pulsar-mongo:27017] INFO org.mongodb.driver.cluster - Monitor thread successfully connected to server with description ServerDescription{addresspulsar-mongo:27017, typeSTANDALONE, stateCONNECTED, oktrue, versionServerVersion{versionList[4, 2, 0]}, minWireVersion0, maxWireVersion8, maxDocumentSize16777216, logicalSessionTimeoutMinutes30, roundTripTimeNanos89058800}这两条日志来自 Mongo Java Driver第一条说明与pulsar-mongo:27017建立了 TCP 连接第二条说明集群监控线程成功连接服务器ServerDescription中给出了服务器类型STANDALONE、状态CONNECTED、oktrue、Mongo 版本4.2.0等关键信息。调试价值如果这里出现连接失败或超时说明mongoUri指向的地址不可达例如容器名解析失败、端口未映射、Mongo 未启动。注意此处主机名是pulsar-mongoDocker 容器名若脱离容器环境运行需要改为实际可达的 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共享订阅多个实例可并行消费receiverQueueSize1000为接收队列大小subscriptionInitialPositionLatest表示新订阅默认从最新消息开始消费acknowledgementsGroupTimeMicros100000表示 100ms 的确认聚合窗口negativeAckRedeliveryDelayMicros60000000表示负确认后 60s 重投递。属性中applicationpulsar-sink、idpublic/default/pulsar-mongo-sink、instance_id0会同步体现在 topic stats 的消费者元数据里。Client 配置serviceUrlpulsar://localhost:6650是 localrun 模式连接 Broker 的服务地址numIoThreads1、numListenerThreads1为 IO/监听线程数connectionsPerBroker1operationTimeoutMs30000keepAliveIntervalSeconds30未开启 TLSuseTlsfalse。调试价值若连接器启动成功但收不到消息优先核对此处serviceUrl是否正确、subscriptionInitialPosition是否为预期位置Latest会跳过此前已发布的历史消息。需要消费历史消息时可通过--subs-position类参数调整订阅初始位置。在集群cluster模式下调试cluster 模式下连接器被提交到 Functions worker 运行调试手段分为两类连接器日志与Admin CLI。使用连接器日志集群模式下一个 worker 上可能同时运行多个连接器实例。要定位某个连接器的日志路径需要使用workerIdworker 实例的唯一标识来确定该连接器运行在哪个 worker 节点上再前往该节点的logs/functions/tenant/namespace/function-name/目录查看对应实例的日志文件。workerId的获取方式见下文status命令。使用 Admin CLIPulsar Admin CLI 提供了三个与调试直接相关的子命令get、status、topics 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 1get获取连接器基础信息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.MongoSinkconfigs配置文件中所有参数batchSize、batchTimeMs以数值形式返回parallelism实例并行度此处为 1processingGuarantees处理保证级别ATLEAST_ONCE表示至少一次语义MongoSink内部采用批量写入 成功/失败分别ack/fail的机制与这一语义一致retainOrdering是否保持消息顺序false表示不做顺序保证autoAck是否自动确认true表示由框架自动完成确认。若get返回的configs与配置文件不一致说明创建命令传入的参数有问题应重新核对--sink-config-file指向的文件内容。status获取连接器运行状态status命令返回实例数量、运行中实例数、每个实例的instanceId、workerId以及各类错误/统计计数是判断连接器是否健康运行的核心命令。更多选项见 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/latestSinkExceptionsSink 写外部系统时的异常计数与最近异常明细如 Mongo 写入失败numReadFromPulsar/numWrittenToSink从 Pulsar 读取的消息数与成功写入 Sink 的消息数两者长期为 0 或严重不对称时都值得深挖lastReceivedTime最近一次收到消息的时间戳为 0 表示从未收到过消息workerId实例所在 worker 的唯一标识。若多个连接器运行在同一 worker 上workerId可帮你定位该连接器运行在哪个节点从而找到对应的日志文件。调试价值latestSystemExceptions与latestSinkExceptions是最直接的故障线索来源。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为空且msgRateIn0说明还没有生产者向该 topic 发布消息订阅public/default/pulsar-mongo-sinkmsgBacklog积压消息数、unackedMessages未确认消息数、msgRateRedeliver重投递速率、typeShared订阅类型。msgBacklog持续增长而msgRateOut0说明消费者没有有效消费应回到status检查 Sink 异常unackedMessages长期不降且blockedSubscriptionOnUnackedMsgstrue说明存在确认堆积问题消费者明细availablePermits可用许可数接近receiverQueueSize1000表示消费者空闲、metadata中的applicationpulsar-sink与idpublic/default/pulsar-mongo-sink可确认该消费者正是本连接器的实例、connectedSince连接建立时间、clientVersion、address消费者地址可用于与status中的workerId对应。调试价值topics stats能够把消息流这一环节单独隔离出来验证——如果 topic 侧一切正常有生产、有消费、无积压、无异常但外部系统里没有数据那么问题很可能出在 Sink 实现或外部系统一侧需要结合连接器日志与外部系统日志继续排查。连接器调试检查清单Checklist以下清单覆盖连接器调试时需要核对的全部关键区域既可作为全面排查的提醒也可作为评估连接器当前状态的工具Pulsar 是否正常启动检查 Broker、BookKeeper、ZooKeeper 以及 standalone 模式下的 Web 服务是否可用bin/pulsar-admin brokers healthcheck类命令可用于快速验证。外部服务是否正常运行确认目标系统本例为 Mongo进程存活、端口可连通、凭据与权限正确必要时在外部系统侧查看其自身日志。nar 包是否完整检查 nar 文件是否下载完整、版本与 Pulsar 版本匹配解压目录META-INF/bundled-dependencies/中依赖是否齐全缺依赖通常表现为class cannot be found。连接器配置文件是否正确对照MongoConfig的必填项与取值范围逐项核对mongoUri、database、collection、batchSize、batchTimeMs并结合get命令返回的configs做双重确认。localrun 模式运行连接器并检查控制台打印的连接器日志nar 解压路径、InstanceConfig、Mongo 连接、Consumer/Client 配置四段信息逐一核对。cluster 模式用get命令获取基础配置信息用status命令获取实例运行状态、异常计数与workerId用topics stats命令获取 topic 及其生产/消费者的统计信息根据workerId定位 worker 节点检查连接器日志文件。进入外部系统验证结果登录 Mongo检查pulsar库messages集合中是否出现符合预期的文档从而确认端到端链路Pulsar → Sink → Mongo是否真正打通。从源码理解调试的关键机制localrun 的两类运行时LocalRunnerpulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java根据配置选择线程模式或进程模式线程模式startThreadedMode→ThreadRuntimeFactory连接器实例运行在pulsar-admin进程内部的独立线程组中日志直接复用当前进程的日志框架因此会实时打印在控制台同时也由框架写入日志文件进程模式startProcessMode→ProcessRuntimeFactory连接器实例以独立子进程运行调试时可通过进程输出重定向或日志文件观察其行为。无论哪种模式每个实例的InstanceConfig都会记录functionId、functionVersion、instanceId、clusterNamelocal、maxBufferedTuples1024等信息这些正是启动日志第二段所展示的内容也是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()从而保证至少一次的处理语义。调试启发batchSize2、batchTimeMs500的配置意味着每凑齐 2 条消息或每 500ms 就会批量写入一次。如果status中numWrittenToSink小于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),仅供参考
