ARTICLE DETAIL

资讯详情

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

Apache Pulsar Go 客户端完整实战指南:生产者、消费者与读者开发详解

Apache Pulsar Go 客户端完整实战指南:生产者、消费者与读者开发详解 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 官方为 GoGolang开发者提供了与 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 包macOSbrew 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 进行依赖管理其中vpulsar:version是版本占位符实际使用时应替换为具体版本号$ dep ensure -add github.com/apache/pulsar/pulsar-client-go/pulsarvpulsar: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 使用pulsarsslscheme 并切换到 TLS 默认端口6651pulsarssl://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 配置参数参数描述默认值URLPulsar 集群的连接 URL无IOThreads处理与 Pulsar broker 连接的线程数1OperationTimeoutSeconds某些 Go 客户端操作创建生产者、订阅/取消订阅 topic的超时时间。重试会一直持续到该阈值之后操作失败30MessageListenerThreads消息监听器消费者和读者使用的线程数1ConcurrentLookupRequests每条 broker 连接上可并发发送的 lookup 请求数。设置上限可避免 broker 过载。只有当客户端需要生产/订阅数千个 topic 时才应调高默认值5000Logger客户端的自定义日志实现一个接收日志级别、文件路径、行号和消息的函数。所有 info/warn/error 消息都会路由到该函数nilTLSTrustCertsFilePath受信任 TLS 证书的文件路径无TLSAllowInsecureConnection客户端是否接受来自 broker 的不可信 TLS 证书falseAuthentication配置认证提供者默认无认证。示例Authentication: NewAuthenticationTLS(my-cert.pem, my-key.pem)nilStatsIntervalInSeconds客户端统计信息发布的间隔秒60实现级补充IOThreads、OperationTimeoutSeconds、ConcurrentLookupRequests、StatsIntervalInSeconds等参数在 C 客户端中都能找到对应实现。例如 ClientConfiguration.h 中setIOThreads、setOperationTimeoutSeconds、setConcurrentLookupRequest、setStatsIntervalInSeconds的 javadoc 明确说明默认值分别为 1、30 秒、50000Go 侧文档取 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()获取生产者的 topicstringName()获取生产者名称stringSend(context.Context, ProducerMessage) error向生产者的 topic 发布消息。该调用会阻塞直到消息被 Pulsar broker 成功确认若超过生产者配置中的SendTimeout则会抛出错误errorSendAsync(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设为true1000MaxPendingMessagesAcrossPartitions跨全部分区的最大待处理消息数。当总量超过该值时会按比例调低单个分区的待处理队列上限对应 C 客户端setMaxPendingMessagesAcrossPartitions默认 50000见 ProducerConfiguration.h50000BlockIfQueueFull设为true时当发送队列满时Send/SendAsync会阻塞而非失败设为false默认时队列满时上述操作会失败并抛出ProducerQueueIsFullErrorfalseMessageRoutingMode消息路由逻辑用于分区 topic。仅当消息未设置 key 时生效。可选round robinpulsar.RoundRobinDistribution默认、全部消息发布到单一分区pulsar.UseSinglePartition、自定义分区方案pulsar.CustomPartitionpulsar.RoundRobinDistributionHashingScheme决定消息发布到哪个分区的哈希函数仅分区 topic。可选pulsar.JavaStringHash等价于 Java 的String.hashCode()、pulsar.Murmur3_32HashMurmur3 哈希、pulsar.BoostHashC Boost 库的哈希函数pulsar.JavaStringHashCompressionType生产者使用的消息数据压缩类型。可选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与setBlockIfQueueFullC 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()返回消费者的 topicstringSubscription()返回消费者的订阅名称stringUnsubcribe()将消费者从指定 topic 退订。若退订操作不成功则抛错errorReceive(context.Context)从 topic 接收单条消息。该方法会阻塞直到有可用消息(Message, error)Ack(Message)向 Pulsar broker 确认一条消息errorAckID(MessageID)按消息 ID 向 broker 确认一条消息errorAckCumulative(Message)确认消息流中一直到并包括指定消息在内的所有消息。AckCumulative会阻塞直到 ack 发送到 broker。此后这些消息将不会被重新投递给消费者。累积确认只能用于 shared 订阅类型errorNack(Message)确认单条消息处理失败errorNackID(MessageID)确认单条消息处理失败按消息 IDerrorClose()关闭消费者使其无法再从 broker 接收消息errorRedeliverUnackedMessages()重新投递 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消息确认超时时间0NackRedeliveryDelay处理失败的消息重新投递前的延迟参见Consumer.Nack()1 分钟SubscriptionType可选值Exclusive、Shared、FailoverExclusiveMessageChannel消费者使用的 Go channel。从 Pulsar topic 到达的消息会传给该 channel无ReceiverQueueSize消费者接收队列大小即应用调用Receive前消费者可累积的消息数。高于默认值 1000 可提升消费者吞吐但会占用更多内存1000MaxTotalReceiverQueueSizeAcrossPartitions设置跨分区的最大接收队列总量。若总量超过该值会降低单个分区的接收队列大小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()返回读者的 topicstringNext(context.Context)接收 topic 上的下一条消息类似于消费者的Receive方法。该方法会阻塞直到有可用消息(Message, error)Close()关闭读者使其无法再从 broker 接收消息errorNext 示例使用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.EarliestMessagetopic 上最早可用消息、pulsar.LatestMessagetopic 上最新可用消息、或用于指定非最早/最新位置的MessageID对象pulsar.LatestMessageMessageChannel读者使用的 Go channel。从 Pulsar topic 到达的消息会传给该 channel无ReceiverQueueSize读者接收队列大小即应用调用Next前读者可累积的消息数。高于默认值 1000 可提升读者吞吐但会占用更多内存1000SubscriptionRolePrefix订阅角色前缀reader消息MessagesPulsar 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 加密传输需要按以下三步配置客户端使用pulsarsslURL 类型将TLSTrustCertsFilePath设置为客户端与 Pulsar broker 使用的 TLS 证书路径配置Authentication选项。完整示例opts : pulsar.ClientOptions{ URL: pulsarssl://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.mdC 客户端安装细节RPM/Deb/Homebrewclient-libraries-cpp.md消息模型、订阅类型Exclusive/Shared/Failover与确认语义concepts-messaging.md分区 topic 与路由概念concepts-architecture-overview.mdTLS 传输加密与认证配置security-tls-transport.md、security-tls-authentication.md消息去重与发送超时设置建议cookbooks-deduplication.mdGo 客户端底层所依赖的 C 客户端源码pulsar-client-cpp赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Go 客户端pulsar-client-go完整使用指南安装、生产者、消费者与 ReaderApache 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消息队列后端流处理上一篇终极黑苹果配置指南OpCore Simplify智能图形化工具完整解析下一篇3步搭建专业在线考试系统学之思开源项目完整部署教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表