How to debug Pulsar connectors
本指南解释如何在本地运行或集群模式中调试连接器,并提供调试检查列表。 为了更好地演示如何调试 Pulsar 连接器,在这里以Mongo sink 连接器作为一个例子。
部署一个 Mongo sink 环境
启动一个 Mongo 服务。
docker pull mongo:4
docker run -d -p 27017:27017 --name pulsar-mongo -v $PWD/data:/data/db mongo:4
创建数据库和集合(Collection)。
docker exec -it pulsar-mongo /bin/bash
mongo
> use pulsar
> db.createCollection('messages')
> exit
启动 Pulsar 单机模式。
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
使用
mongo-sink-config.yaml
文件配置 Mongo sink。configs:
mongoUri: "mongodb://pulsar-mongo:27017"
database: "pulsar"
collection: "messages"
batchSize: 2
batchTimeMs: 500
docker cp mongo-sink-config.yaml pulsar-mongo-standalone:/pulsar/
下载 Mongo sink 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
在本地运行模式下调试
使用 localrun
命令以本地模式启动 Mongo sink。
Tip
关于
localrun
命令的更多信息,请参阅localrun
。
./bin/pulsar-admin sinks localrun \
--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
使用连接器日志
Use one of the following methods to get a connector log in localrun mode:
After executing the
localrun
command, the log is automatically printed on the console.The log is located at:
logs/functions/tenant/namespace/function-name/function-name-instance-id.log
示例
The path of the Mongo sink connector is:
logs/functions/public/default/pulsar-mongo-sink/pulsar-mongo-sink-0.log
为了清楚地解释日志信息,此处将大块信息分解成小块,并为每个块添加描述。
此日志信息显示解压后 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/,
Tip
If
class cannot be found
exception is thrown, check whether the nar file is decompressed in the folderfile:/tmp/pulsar-nar/pulsar-io-mongo-2.4.0.nar-unpacked/META-INF/bundled-dependencies/
or not.This piece of log information illustrates the basic information about the Mongo sink connector, such as tenant, namespace, name, parallelism, resources, and so on, which can be used to check whether the Mongo sink connector is configured correctly or not.
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)
此日志信息显示与 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}
该日志信息说明了消费者和客户端配置,包括主题名称、订阅名称、订阅类型等等。
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
}
在集群模式中调试
You can use the following methods to debug a connector in cluster mode:
使用连接器日志
在集群模式下,多个连接器可以运行在一个 worker 上。 要找到指定连接器的日志路径,请使用 workerId
来定位连接器日志。
使用管理命令行工具(admin CLI)
Pulsar admin CLI helps you debug Pulsar connectors with the following subcommands:
创建 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
命令,获取 Mongo sink 连接器的基本信息,例如租户、命名空间、名称、并行度等等。
./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
}
Tip
更多
get
命令信息,请参阅get
章节。
status
使用 status
命令获取 Mongo sink 连接器的当前状态,例如实例数量、正在运行实例的数量、instanceId、workerId 等。
./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"
}
} ]
}
Tip
关于
status
命令的更多信息,可参阅status
。如果一个 worker 上运行着多个连接器,
workerId
可以找到指定的连接器所运行的 worker。
topics stats
使用 topics stats
命令获取主题及其关联的生产者、消费者的统计信息,如主题是否已收到消息,或者是否有消息积压,或可用权限以及其它关键信息。 All rates are computed over a 1-minute window and are relative to the last completed 1-minute period.
./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"
}
Tip
关于
topic stats
命令的更多信息,请参阅topic stats
。
检查清单
此清单列出了调试连接器时要检查的主要事项。 该清单提醒我们应该注意什么,以确保彻底检查,并作为评估工具获取连接器状态。
Pulsar 是否成功启动?
外部服务运行正常吗?
nar 包是否完整?
连接器配置文件正确吗?
在本地运行模式下,运行连接器并检查控制台输出的信息(连接器日志)。
在集群模式中:
使用
get
命令获取基本信息。使用
status
命令获取当前状态。使用
topics stats
命令获取特定主题及其关联的生产者和消费者的统计信息。检查连接器日志。
进入外部系统并验证结果。