- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Apache Pulsar 官方为 Go(Golang)开发者提供了与 Java、C++ 客户端能力对齐的 Go 客户端库,用于创建消息生产者(Producer)、消费者(Consumer)与读者(Reader)。本指南以 Pulsar 2.3.1 版本文档为骨架,结合本仓库中 C++ 客户端源码(Go 客户端基于 C++ 客户端库封装),完整覆盖安装方式、连接 URL 格式、Client/Producer/Consumer/Reader 的全部配置参数、消息构造与 TLS 加密认证配置,并补充实现级原理说明,帮助你从零开始构建可运行的 Pulsar Go 应用。
说明:本仓库当前包含的是 Pulsar 2.3.1 时代 Go 客户端(基于 C++ 客户端库的 CGo 封装)的官方文档(见 client-libraries-go.md)。在后续版本中,Apache Pulsar 推出了独立的纯 Go 客户端(
github.com/apache/pulsar-client-go)并逐步弃用 CGo 封装,本指南严格以 2.3.1 版本文档内容为准。
安装与依赖
前置要求
Pulsar Go 客户端库基于 C++ 客户端库实现,因此在安装 Go 包之前,需要先按照 C++ 客户端的安装指引安装二进制库。安装方式有三种(详见 C++ 客户端文档):
- RPM 包:适用于基于 RPM 的 Linux 发行版;
- Deb 包:适用于 Debian/Ubuntu 系发行版;
- Homebrew 包(macOS):
brew install libpulsar一类的方式安装。
兼容性警告:Go 客户端的版本号必须与Pulsar C++ 客户端库的版本号严格一致。C++ 客户端库正是本仓库中 pulsar-client-cpp 目录对应的模块。
安装 Go 包
安装pulsarGo 库可以使用go get:
$ go get -u github.com/apache/pulsar/pulsar-client-go/pulsar注意:
go get不支持拉取指定 tag,它总是拉取 master 分支版本的 Go 客户端,因此你需要一个与 master 匹配的 C++ 客户端库。
也可以使用 dep 进行依赖管理(其中@v@pulsar:version@是版本占位符,实际使用时应替换为具体版本号):
$ dep ensure -add github.com/apache/pulsar/pulsar-client-go/pulsar@v@pulsar:version@安装完成后即可在项目中导入:
import "github.com/apache/pulsar/pulsar-client-go/pulsar"连接 URL
使用任何 Pulsar 客户端库连接集群,都需要指定一个 Pulsar 二进制协议 URL。Pulsar 协议 URL 归属于特定集群,使用pulsarscheme,默认端口为6650。本机示例:
pulsar://localhost:6650生产环境集群的 URL 形如:
pulsar://pulsar.us-west.example.com:6650如果使用 TLS 认证,URL 使用pulsar+sslscheme 并切换到 TLS 默认端口6651:
pulsar+ssl://pulsar.us-west.example.com:6651创建客户端(Client)
要与 Pulsar 交互,第一步是创建一个Client对象。使用NewClient函数并传入一个ClientOptions对象即可:
import ( "log" "runtime" "github.com/apache/pulsar/pulsar-client-go/pulsar" ) func main() { client, err := pulsar.NewClient(pulsar.ClientOptions{ URL: "pulsar://localhost:6650", OperationTimeoutSeconds: 5, MessageListenerThreads: runtime.NumCPU(), }) if err != nil { log.Fatalf("Could not instantiate Pulsar client: %v", err) } }Client 配置参数
| 参数 | 描述 | 默认值 |
|---|---|---|
URL | Pulsar 集群的连接 URL | 无 |
IOThreads | 处理与 Pulsar broker 连接的线程数 | 1 |
OperationTimeoutSeconds | 某些 Go 客户端操作(创建生产者、订阅/取消订阅 topic)的超时时间。重试会一直持续到该阈值,之后操作失败 | 30 |
MessageListenerThreads | 消息监听器(消费者和读者)使用的线程数 | 1 |
ConcurrentLookupRequests | 每条 broker 连接上可并发发送的 lookup 请求数。设置上限可避免 broker 过载。只有当客户端需要生产/订阅数千个 topic 时才应调高默认值 | 5000 |
Logger | 客户端的自定义日志实现(一个接收日志级别、文件路径、行号和消息的函数)。所有 info/warn/error 消息都会路由到该函数 | nil |
TLSTrustCertsFilePath | 受信任 TLS 证书的文件路径 | 无 |
TLSAllowInsecureConnection | 客户端是否接受来自 broker 的不可信 TLS 证书 | false |
Authentication | 配置认证提供者(默认无认证)。示例:Authentication: NewAuthenticationTLS("my-cert.pem", "my-key.pem") | nil |
StatsIntervalInSeconds | 客户端统计信息发布的间隔(秒) | 60 |
实现级补充:IOThreads、OperationTimeoutSeconds、ConcurrentLookupRequests、StatsIntervalInSeconds等参数在 C++ 客户端中都能找到对应实现。例如 ClientConfiguration.h 中setIOThreads、setOperationTimeoutSeconds、setConcurrentLookupRequest、setStatsIntervalInSeconds的 javadoc 明确说明默认值分别为 1、30 秒、50000(Go 侧文档取 5000 为更低的安全上限)和 600 秒;其中setConcurrentLookupRequest明确指出该值"仅在需要基于单个客户端生产/订阅数千个 topic 时才应调高"。TLSAllowInsecureConnection与TLSTrustCertsFilePath分别对应 C++ 侧setTlsAllowInsecureConnection与setTlsTrustCertsFilePath,后者描述为"设置受信任 TLS 证书文件的路径"。
生产者(Producers)
生产者负责向 Pulsar topic 发布消息。使用ProducerOptions对象配置 Go 生产者:
producer, err := client.CreateProducer(pulsar.ProducerOptions{ Topic: "my-topic", }) if err != nil { log.Fatalf("Could not instantiate Pulsar producer: %v", err) } defer producer.Close() msg := pulsar.ProducerMessage{ Payload: []byte("Hello, Pulsar"), } if err := producer.Send(msg); err != nil { log.Fatalf("Producer could not send message: %v", err) }阻塞操作:创建新的 Pulsar 生产者时,操作会阻塞(等待一个 go channel),直到生产者创建成功或抛出错误。
生产者操作(Producer operations)
| 方法 | 描述 | 返回类型 |
|---|---|---|
Topic() | 获取生产者的 topic | string |
Name() | 获取生产者名称 | string |
Send(context.Context, ProducerMessage) error | 向生产者的 topic 发布消息。该调用会阻塞直到消息被 Pulsar broker 成功确认;若超过生产者配置中的SendTimeout则会抛出错误 | error |
SendAsync(context.Context, ProducerMessage, func(ProducerMessage, error)) | 异步向生产者的 topic 发布消息。第三个参数是回调函数,在消息被确认或抛出错误时执行 | 无 |
Close() | 关闭生产者并释放其所有资源。Close()被调用后将不再接受新消息。该方法会阻塞直到所有待发送的发布请求被 Pulsar 持久化。若抛出错误,则不再重试任何待写入的消息 | error |
下面是一个更完整的生产者示例,同时演示同步与异步发送:
import ( "context" "fmt" "log" "github.com/apache/pulsar/pulsar-client-go/pulsar" ) func main() { // 实例化 Pulsar 客户端 client, err := pulsar.NewClient(pulsar.ClientOptions{ URL: "pulsar://localhost:6650", }) if err != nil { log.Fatal(err) } // 使用客户端实例化生产者 producer, err := client.CreateProducer(pulsar.ProducerOptions{ Topic: "my-topic", }) if err != nil { log.Fatal(err) } ctx := context.Background() // 同步发送 10 条消息、异步发送 10 条消息 for i := 0; i < 10; i++ { // 创建一条消息 msg := pulsar.ProducerMessage{ Payload: []byte(fmt.Sprintf("message-%d", i)), } // 尝试同步发送消息 if err := producer.Send(ctx, msg); err != nil { log.Fatal(err) } // 创建另一条用于异步发送的消息 asyncMsg := pulsar.ProducerMessage{ Payload: []byte(fmt.Sprintf("async-message-%d", i)), } // 尝试异步发送消息并处理回调 producer.SendAsync(ctx, asyncMsg, func(msg pulsar.ProducerMessage, err error) { if err != nil { log.Fatal(err) } fmt.Printf("Message %s successfully published", msg.ID()) }) } }生产者配置(Producer configuration)
| 参数 | 描述 | 默认值 |
|---|---|---|
Topic | 生产者将发布消息的 Pulsar topic | 无 |
Name | 生产者名称。若未显式指定,Pulsar 会自动生成一个全局唯一名称,之后可用Name()方法获取。若显式指定名称,则必须跨所有 Pulsar 集群唯一,否则创建操作会抛错 | 自动生成 |
SendTimeout | 向 topic 发布消息时,生产者会等待负责的 Pulsar broker 确认。若消息在SendTimeout阈值内未被确认,则抛出错误。若设为 -1,超时设为无穷大(即移除超时)。使用 Pulsar 消息去重 功能时建议移除发送超时 | 30 秒 |
MaxPendingMessages | 待处理消息队列(即等待 broker 确认的消息)的最大大小。默认情况下,队列满时所有Send和SendAsync调用都会失败,除非BlockIfQueueFull设为true | 1000 |
MaxPendingMessagesAcrossPartitions | 跨全部分区的最大待处理消息数。当总量超过该值时,会按比例调低单个分区的待处理队列上限(对应 C++ 客户端setMaxPendingMessagesAcrossPartitions,默认 50000,见 ProducerConfiguration.h) | 50000 |
BlockIfQueueFull | 设为true时,当发送队列满时Send/SendAsync会阻塞而非失败;设为false(默认)时,队列满时上述操作会失败并抛出ProducerQueueIsFullError | false |
MessageRoutingMode | 消息路由逻辑(用于分区 topic)。仅当消息未设置 key 时生效。可选:round robin(pulsar.RoundRobinDistribution,默认)、全部消息发布到单一分区(pulsar.UseSinglePartition)、自定义分区方案(pulsar.CustomPartition) | pulsar.RoundRobinDistribution |
HashingScheme | 决定消息发布到哪个分区的哈希函数(仅分区 topic)。可选:pulsar.JavaStringHash(等价于 Java 的String.hashCode())、pulsar.Murmur3_32Hash(Murmur3 哈希)、pulsar.BoostHash(C++ Boost 库的哈希函数) | pulsar.JavaStringHash |
CompressionType | 生产者使用的消息数据压缩类型。可选:LZ4、ZLIB、ZSTD | 不压缩 |
MessageRouter | 默认情况下,Pulsar 对分区 topic 使用 round-robin 路由。MessageRouter允许通过一个接收 Pulsar 消息和 topic 元数据、返回整数的函数指定自定义路由逻辑,函数签名为func(Message, TopicMetadata) int | 无 |
实现级补充:
- 哈希方案:三种
HashingScheme在 C++ 客户端中分别有独立实现。JavaStringHash.cc 用hash = 31 * hash + val[i]的经典 JavaString.hashCode()算法计算并对int32取正;Murmur3_32Hash与BoostHash分别位于 Murmur3_32Hash.cc 与 BoostHash.cc,三个类都实现了同一个makeHash接口。 - 路由实现:默认的 round-robin 路由在 RoundRobinMessageRouter.cc 中实现:当消息带 key 时直接对 key 哈希取模选分区;无 key 且未开启批量时按消息粒度轮询;开启批量时则在同一分区累积消息,直到消息数、批量字节数或最大批量延迟任一阈值被触发才切换分区。这解释了
MessageRoutingMode与HashingScheme协作时的真实行为。 - 队列控制:
MaxPendingMessages与BlockIfQueueFull对应 C++ 侧setMaxPendingMessages(默认 1000)与setBlockIfQueueFull;C++ javadoc 明确写道"当队列满时,默认情况下所有Producer::send与Producer::sendAsync调用都会失败,除非blockIfQueueFull为 true",与 Go 文档描述一致(见 ProducerConfiguration.h)。 - 压缩类型:C++ 侧
setCompressionType的 javadoc 补充了重要限制:ZSTD 自 Pulsar 2.3 起支持,但要求消费端应用版本也必须 >= 2.3(见 ProducerConfiguration.h)。
消费者(Consumers)
消费者订阅一个或多个 Pulsar topic 并监听其上产生的消息。使用ConsumerOptions对象配置 Go 消费者。下面是一个使用 channel 的基础示例:
msgChannel := make(chan pulsar.ConsumerMessage) consumerOpts := pulsar.ConsumerOptions{ Topic: "my-topic", SubscriptionName: "my-subscription-1", Type: pulsar.Exclusive, MessageChannel: msgChannel, } consumer, err := client.Subscribe(consumerOpts) if err != nil { log.Fatalf("Could not establish subscription: %v", err) } defer consumer.Close() for cm := range msgChannel { msg := cm.Message fmt.Printf("Message ID: %s", msg.ID()) fmt.Printf("Message value: %s", string(msg.Payload())) consumer.Ack(msg) }阻塞操作:创建新的 Pulsar 消费者时,操作会阻塞(在一个 go channel 上),直到订阅成功建立或抛出错误。
消费者操作(Consumer operations)
| 方法 | 描述 | 返回类型 |
|---|---|---|
Topic() | 返回消费者的 topic | string |
Subscription() | 返回消费者的订阅名称 | string |
Unsubcribe() | 将消费者从指定 topic 退订。若退订操作不成功则抛错 | error |
Receive(context.Context) | 从 topic 接收单条消息。该方法会阻塞直到有可用消息 | (Message, error) |
Ack(Message) | 向 Pulsar broker 确认一条消息 | error |
AckID(MessageID) | 按消息 ID 向 broker 确认一条消息 | error |
AckCumulative(Message) | 确认消息流中一直到并包括指定消息在内的所有消息。AckCumulative会阻塞直到 ack 发送到 broker。此后这些消息将不会被重新投递给消费者。累积确认只能用于 shared 订阅类型 | error |
Nack(Message) | 确认单条消息处理失败 | error |
NackID(MessageID) | 确认单条消息处理失败(按消息 ID) | error |
Close() | 关闭消费者,使其无法再从 broker 接收消息 | error |
RedeliverUnackedMessages() | 重新投递 topic 上所有未确认消息。在 failover 模式下,若消费者不是该 topic 上的活跃消费者则该请求被忽略;在 shared 模式下,重新投递的消息会分布到连接该 topic 的所有消费者。注意:这是一个非阻塞操作,不抛错 | 无 |
Receive 示例
使用Receive()方法处理消息的消费者示例:
import ( "context" "log" "github.com/apache/pulsar/pulsar-client-go/pulsar" ) func main() { // 实例化 Pulsar 客户端 client, err := pulsar.NewClient(pulsar.ClientOptions{ URL: "pulsar://localhost:6650", }) if err != nil { log.Fatal(err) } // 使用客户端对象实例化消费者 consumer, err := client.Subscribe(pulsar.ConsumerOptions{ Topic: "my-golang-topic", SubscriptionName: "sub-1", SubscriptionType: pulsar.Exclusive, }) if err != nil { log.Fatal(err) } defer consumer.Close() ctx := context.Background() // 在 topic 上无限监听 for { msg, err := consumer.Receive(ctx) if err != nil { log.Fatal(err) } // 对消息做处理 err = processMessage(msg) if err == nil { // 消息处理成功 consumer.Ack(msg) } else { // 消息处理失败 consumer.Nack(msg) } } }消费者配置(Consumer configuration)
| 参数 | 描述 | 默认值 |
|---|---|---|
Topic | 消费者将建立订阅并监听消息的 Pulsar topic | 无 |
SubscriptionName | 该消费者的订阅名称 | 无 |
Name | 消费者名称 | 无 |
AckTimeout | 消息确认超时时间 | 0 |
NackRedeliveryDelay | 处理失败的消息重新投递前的延迟(参见Consumer.Nack()) | 1 分钟 |
SubscriptionType | 可选值:Exclusive、Shared、Failover | Exclusive |
MessageChannel | 消费者使用的 Go channel。从 Pulsar topic 到达的消息会传给该 channel | 无 |
ReceiverQueueSize | 消费者接收队列大小,即应用调用Receive前消费者可累积的消息数。高于默认值 1000 可提升消费者吞吐,但会占用更多内存 | 1000 |
MaxTotalReceiverQueueSizeAcrossPartitions | 设置跨分区的最大接收队列总量。若总量超过该值,会降低单个分区的接收队列大小 | 50000 |
读者(Readers)
读者从 Pulsar topic 处理消息。与消费者的区别在于:使用读者时必须显式指定要从流中的哪条消息开始处理(消费者则自动从最近的未确认消息开始)。使用ReaderOptions对象配置 Go 读者:
reader, err := client.CreateReader(pulsar.ReaderOptions{ Topic: "my-golang-topic", StartMessageId: pulsar.LatestMessage, })阻塞操作:创建新的 Pulsar 读者时,操作会阻塞(在一个 go channel 上),直到读者成功创建或抛出错误。
读者操作(Reader operations)
| 方法 | 描述 | 返回类型 |
|---|---|---|
Topic() | 返回读者的 topic | string |
Next(context.Context) | 接收 topic 上的下一条消息(类似于消费者的Receive方法)。该方法会阻塞直到有可用消息 | (Message, error) |
Close() | 关闭读者,使其无法再从 broker 接收消息 | error |
"Next" 示例
使用Next()方法处理消息的读者示例:
import ( "context" "log" "github.com/apache/pulsar/pulsar-client-go/pulsar" ) func main() { // 实例化 Pulsar 客户端 client, err := pulsar.NewClient(pulsar.ClientOptions{ URL: "pulsar://localhost:6650", }) if err != nil { log.Fatalf("Could not create client: %v", err) } // 使用客户端实例化读者 reader, err := client.CreateReader(pulsar.ReaderOptions{ Topic: "my-golang-topic", StartMessageID: pulsar.EarliestMessage, }) if err != nil { log.Fatalf("Could not create reader: %v", err) } defer reader.Close() ctx := context.Background() // 在 topic 上监听消息 for { msg, err := reader.Next(ctx) if err != nil { log.Fatalf("Error reading from topic: %v", err) } // 处理消息 } }在上面的示例中,读者从最早可用消息(pulsar.EarliestMessage)开始读取。读者也可以从最新消息(pulsar.LatestMessage)开始,或通过DeserializeMessageID函数从字节数组反序列化得到某个指定的MessageID作为起点:
lastSavedId := // Read last saved message id from external store as byte[] reader, err := client.CreateReader(pulsar.ReaderOptions{ Topic: "my-golang-topic", StartMessageID: DeserializeMessageID(lastSavedId), })这种从外部存储恢复MessageID的用法,使读者天然适合"断点续读"场景:例如将上次处理到的消息 ID 序列化后持久化,重启后精确地从断点继续消费。
读者配置(Reader configuration)
| 参数 | 描述 | 默认值 |
|---|---|---|
Topic | 读者将建立订阅并监听消息的 Pulsar topic | 无 |
Name | 读者名称 | 无 |
StartMessageID | 初始读者位置,即读者开始处理消息的位置。可选:pulsar.EarliestMessage(topic 上最早可用消息)、pulsar.LatestMessage(topic 上最新可用消息)、或用于指定非最早/最新位置的MessageID对象 | pulsar.LatestMessage |
MessageChannel | 读者使用的 Go channel。从 Pulsar topic 到达的消息会传给该 channel | 无 |
ReceiverQueueSize | 读者接收队列大小,即应用调用Next前读者可累积的消息数。高于默认值 1000 可提升读者吞吐,但会占用更多内存 | 1000 |
SubscriptionRolePrefix | 订阅角色前缀 | reader |
消息(Messages)
Pulsar Go 客户端提供了ProducerMessage接口,用于构造要发布到 Pulsar topic 的消息。示例:
msg := pulsar.ProducerMessage{ Payload: []byte("Here is some message data"), Key: "message-key", Properties: map[string]string{ "foo": "bar", }, EventTime: time.Now(), ReplicationClusters: []string{"cluster1", "cluster3"}, } if err := producer.send(msg); err != nil { log.Fatalf("Could not publish message due to: %v", err) }ProducerMessage对象可用参数如下:
| 参数 | 描述 |
|---|---|
Payload | 消息的实际数据载荷 |
Key | 消息关联的可选 key(对 topic 压缩等场景尤其有用) |
Properties | 附加到消息上的应用自定义元数据键值对(key 和 value 都必须是字符串) |
EventTime | 与消息关联的时间戳 |
ReplicationClusters | 消息将被复制到的集群列表。Pulsar broker 自动处理消息复制;仅当需要覆盖 broker 默认设置时才应修改此项 |
TLS 加密与认证
要启用 TLS 加密传输,需要按以下三步配置客户端:
- 使用
pulsar+sslURL 类型; - 将
TLSTrustCertsFilePath设置为客户端与 Pulsar broker 使用的 TLS 证书路径; - 配置
Authentication选项。
完整示例:
opts := pulsar.ClientOptions{ URL: "pulsar+ssl://my-cluster.com:6651", TLSTrustCertsFilePath: "/path/to/certs/my-cert.csr", Authentication: NewAuthenticationTLS("my-cert.pem", "my-key.pem"), }实现级补充:TLS 相关配置项在 C++ 客户端中同样一一对应:ClientConfiguration.h 提供了setUseTls(默认 false)、setTlsTrustCertsFilePath、setTlsAllowInsecureConnection(默认 false)以及setValidateHostName(默认 false)等方法。其中setValidateHostName启用 TLS 主机名校验,会按 RFC 2818 的服务器身份主机名验证规则,用 x509 证书的 CN/SAN 与期望的 broker 主机名比对——生产环境建议在 TLS 之外同时开启主机名校验以防范中间人攻击。
延伸阅读
- 客户端库总览与功能矩阵:client-libraries.md
- C++ 客户端安装细节(RPM/Deb/Homebrew):client-libraries-cpp.md
- 消息模型、订阅类型(Exclusive/Shared/Failover)与确认语义:concepts-messaging.md
- 分区 topic 与路由概念:concepts-architecture-overview.md
- TLS 传输加密与认证配置:security-tls-transport.md、security-tls-authentication.md
- 消息去重与发送超时设置建议:cookbooks-deduplication.md
- Go 客户端底层所依赖的 C++ 客户端源码:pulsar-client-cpp
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar Go 客户端(pulsar-client-go)完整使用指南:安装、生产者、消费者与 Reader
Apache Pulsar Go 客户端(pulsar client go)完整使用指南:安装、生产者、消费者与 Reader 导读 本文基于 Apache P
消息队列后端流处理Apache Pulsar Go 客户端实战指南:基于 C++ 客户端库的生产者、消费者与 Reader 开发
Apache Pulsar Go 客户端实战指南:基于 C++ 客户端库的生产者、消费者与 Reader 开发 本指南以 Apache Pulsar 2.3.0
消息队列后端流处理Apache Pulsar CGo 客户端(pulsar-client-go)开发指南:安装、配置与生产者/消费者/Reader 实战
Apache Pulsar CGo 客户端(pulsar client go)开发指南:安装、配置与生产者/消费者/Reader 实战 Apache Pulsa
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考