ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Watermill 生态扩展实战指南:Awesome 第三方库清单、日志与可观测性适配器及自定义 Pub/Sub

Watermill 生态扩展实战指南:Awesome 第三方库清单、日志与可观测性适配器及自定义 Pub/Sub Watermill 生态扩展实战指南Awesome 第三方库清单、日志与可观测性适配器及自定义 Pub/Sub【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermillWatermill 的核心设计是官方内核 可插拔适配器框架只定义消息模型与路由器而具体的消息中间件Pub/Sub、日志、追踪全部通过接口接入。官方将社区中未经其维护的优质第三方库集中收录于 Awesome Watermill 清单即本仓库docs/content/docs/awesome.md按示例项目、Pub/Sub 适配器、日志适配器、可观测性与其他工具分类。阅读本文后你将掌握该清单的完整构成、第三方库与 Watermill 的接入接口原理并能依据仓库内的通用测试套件评估一个第三方适配器的成熟度甚至独立实现并贡献你自己的适配器。官方与第三方生态的边界Awesome 清单的定位Awesome 清单的第一段文字就划清了责任边界下列库并非由 Three Dots Labs 维护官方无法提供支持、也不保证它们一定工作正常使用前需要自行调研Do your own research。这正是它被称为 Selected unofficial libraries 的原因——官方收录它们只是因为你可能会觉得有用而不是背书。该清单本身是一个可维护的文档文件如果你知道其他值得收录的库或者你本人就是某个适配器的作者可以通过编辑 docs/content/docs/awesome.md 把它加入清单。整个清单分为五类Examples社区编写的完整示例项目Pub/Subs针对各类消息中间件的第三方适配器Logging将 Watermill 日志接入 logrus / zap / zerolog 的适配器Observability接入 OpenCensus / OpenTelemetry 的追踪与指标适配器Other代码生成、工具集等周边设施。清单末尾还特别引导读者如果想了解如何实现自己的 Pub/Sub 适配器参见 Implementing custom Pub/Sub。理解接入点第三方库如何与 Watermill 衔接要评估任何一个第三方库先要理解 Watermill 为它们定义的插槽。核心接入点全部在 message/pubsub.go 中任何消息中间件适配器本质上都是这两个接口的实现// Publisher is the emitting part of a Pub/Sub. type Publisher interface { // Publish publishes provided messages to the given topic. // Publish can be synchronous or asynchronous - it depends on the implementation. // Most publisher implementations dont support atomic publishing of messages. // Publish does not work with a single Context. Use the Context() method of each message instead. // Publish must be thread safe. Publish(topic string, messages ...*Message) error // Close should flush unsent messages if publisher is async. Close() error } // Subscriber is the consuming part of the Pub/Sub. type Subscriber interface { // Subscribe returns an output channel with messages from the provided topic. // The channel is closed after Close() is called on the subscriber. // To receive the next message, Ack() must be called on the received message. // If message processing fails and the message should be redelivered Nack() should be called instead. Subscribe(ctx context.Context, topic string) (-chan *Message, error) // Close closes all subscriptions with their output channels and flushes offsets etc. when needed. Close() error }从接口注释可以提炼出第三方适配器必须遵守的关键契约Publish必须线程安全且不保证原子性——批量发布多条消息时一旦某条失败后续消息不会继续发布见 docs/content/docs/pub-sub.md 的 Publishing multiple messages 一节同步/异步由实现决定异步发布器必须在Close()时冲刷未发送的消息忘记关闭发布器可能丢失消息Ack/Nack 是 Subscriber 的职责正确实现应在收到上一条消息的Ack/Nack后再消费下一条且必须等 Watermill 处理完消息、返回 Ack 之后才向底层 broker 提交 offset否则进程在消息处理完成前崩溃会丢消息另有可选的 SubscribeInitializer 接口含SubscribeInitialize(topic string) error用于在消费前初始化订阅非强制实现。对日志类第三方库而言接入点则是 log.go 中定义的LoggerAdapter接口type LoggerAdapter interface { Error(msg string, err error, fields LogFields) Info(msg string, fields LogFields) Debug(msg string, fields LogFields) Trace(msg string, fields LogFields) With(fields LogFields) LoggerAdapter }Watermill 自身提供了写入标准输出的 StdLoggerAdapterNewStdLogger(debug, trace bool)与静默的NopLogger第三方日志适配器的工作就是把LoggerAdapter的方法调用桥接到 logrus、zap、zerolog 各自的 API 上。逐类盘点 Awesome 清单中的第三方库以下内容完整覆盖 docs/content/docs/awesome.md 中收录的每一类库。需要说明的是清单原文以外部仓库链接形式给出此处按名称与功能分类整理便于按图索骥。示例项目Examples清单收录了三个社区示例用于学习如何在真实项目中组合使用 Watermillgolang-taipei-watermill-exampleWatermill 的完整示例项目Kafka-PubSub以 Kafka 为消息中间件的示例go-example-financing一个金融场景融资的 Go 示例。这些项目与仓库内官方维护的_examples/目录互补——官方示例覆盖了 基础应用、路由器、CQRS、指标 以及各类 Pub/Sub 接入示例社区示例则展示更多真实业务组合。Pub/Sub 适配器第三方这是清单中体量最大的一类。每个适配器都实现了上文所述的message.Publisher/message.Subscriber接口把 Watermill 统一的消息模型映射到特定中间件。清单收录了以下第三方适配器目标中间件适配器定位AMQP 1.0watermill-amqp10支持 AMQP 1.0 协议区别于官方内置的 AMQP/RabbitMQ 适配器Apache Pulsarwatermill-pulsar云原生流式消息平台 PulsarApache RocketMQwatermill-rocketmq阿里巴巴开源的分布式消息中间件CockroachDBwatermill-crdb基于分布式数据库 CockroachDB 的持久化 Pub/SubEnsignwatermill-ensignRotational 的事件基础设施服务GoogleCloud Pub/Sub HTTP Pushwatermill-googlecloud-http以 HTTP Push 方式接入 GoogleCloud Pub/SubMongoDBwatermill-mongodb以 MongoDB 作为消息存储MQTTwatermill-mqtt物联网场景常用的轻量级消息协议NSQwatermill-nsq分布式实时消息平台 NSQRedis Zsetwatermill-rediszset基于 Redis 有序集合的轻量 Pub/SubSQLitewatermill-comfymill以 SQLite 为后端的单机持久化从这些条目可以看出生态的多样性除了传统消息队列还可以用数据库MongoDB、CockroachDB、SQLite或 Redis 结构实现发布/订阅。官方内置支持的中间件清单见 docs/content/pubsubs/_index.mdKafka、AMQP、GoChannel、Redis Stream、GoogleCloud、NATS、SQL 等第三方清单正是对它的有力补充。日志适配器Logging清单按日志库分三组每组都有多个可选适配器logruswatermill-logrus-adapter与walrus两个适配器zap两个同名watermillzap适配器来自不同作者zerologzerowater、watermillzlog、zerolog-watermill-adapter三个适配器。选择时可以根据团队已有日志栈决定。接入方式统一为构造一个实现 LoggerAdapter 的适配器实例再将其传入需要日志的组件构造函数——例如 gochannel.NewGoChannel 的第二个参数logger watermill.LoggerAdapter传nil时内部自动回退到watermill.NopLogger{}。日志适配器的With(fields LogFields)方法用于携带上下文字段如pubsub_uuid便于关联同一组件实例产生的所有日志。可观测性Observability这一组把 Watermill 的组件接入业界主流观测栈OpenCensuswatermill-opencensus与ocwatermill两个适配器OpenTelemetrywatermill-opentelemetry、watermill-opentelemetry-go-extra、watermill-opentelemetry另一作者以及针对 AMQP 的otel-watermill-amqp和针对 GoChannel 的watermill-otel-tracable-gochannel。值得注意 AMQP 与 GoChannel 都有专门的追踪适配器说明消息链路追踪需要针对传输层实现细节做处理例如把 trace 上下文注入/提取到消息的 Metadata。仓库内官方的观测能力集中在 components/metrics 组件基于 Prometheus 的指标构建器、Handler 装饰器等第三方可观测性库则主要面向分布式追踪两者可以搭配使用。其他工具Othergo-watermill-templateAsyncAPI 社区提供的代码生成模板可从 AsyncAPI 规范生成 Watermill 相关代码watermillx一套 Watermill 扩展工具集protoc-gen-event从 Protobuf 定义生成事件代码的插件与仓库内 CQRS 组件的 Protobuf marshaler 思路一脉相承。如何评估第三方适配器用通用测试套件做体检既然官方声明不保证第三方库可用读者就必须自己验证。Watermill 为此提供了标准工具——通用 Pub/Sub 测试套件位于 pubsub/tests/test_pubsub.go。任何自称生产可用的 Pub/Sub 实现都应能通过它测试套件本身也常被第三方适配器的作者直接复用。入口函数签名如下func TestPubSub( t *testing.T, features Features, pubSubConstructor PubSubConstructor, consumerGroupPubSubConstructor ConsumerGroupPubSubConstructor, )它内部自动运行一批覆盖基础与边界场景的测试TestPublishSubscribe基础收发、TestConcurrentSubscribe50 个并发订阅者、TestConcurrentSubscribeMultipleTopics、TestResendOnErrorNack 后重投递、TestNoAck未 Ack 不推送下一条、TestContinueAfterSubscribeClose、TestConcurrentClose、TestContinueAfterErrors、TestPublishSubscribeInOrder、TestPublisherClose、TestTopic、TestMessageCtx、TestSubscribeCtx、TestNewSubscriberReceivesOldMessages与TestReconnect外加消费者组测试TestConsumerGroups。另有 TestPubSubStressTest 可通过环境变量STRESS_TEST_COUNT控制压力轮数默认 10 轮。Features结构体是评估适配器能力边界的关键也是阅读第三方库文档时的对照表字段含义ConsumerGroups是否支持消费者组ExactlyOnceDelivery是否支持恰好一次投递GuaranteedOrder是否保证消息顺序GuaranteedOrderWithSingleSubscriber仅单个订阅者时是否保证顺序Persistent消息是否持久化GoChannel 不支持RestartServiceCommand用于测试断线重连的重启 broker 的命令RequireSingleInstance是否要求单实例工作如 GoChannelNewSubscriberReceivesOldMessages新订阅者能否收到历史消息如 KafkaContextPreserved是否保留消息发布时的 Context对照仓库内官方适配器的特性表可以直观感受特性矩阵的写法例如 docs/content/pubsubs/kafka.md 标明 Kafka 适配器支持消费者组、持久化与保证顺序需配合分区键但不支持恰好一次投递而 docs/content/pubsubs/gochannel.md 的 GoChannel 恰好相反——不支持消费者组与持久化。评估第三方适配器时也应要求作者提供同样格式的特性表。自行实现并贡献一个适配器如果清单中没有你需要的中间件Implementing custom Pub/Sub 给出了标准路径只需实现message.Publisher与message.Subscriber两个接口然后用通用测试套件验证。结合该文档与源码实现时不要遗漏以下检查清单日志提供清晰、分级的日志消息复用LoggerAdapter的 Error/Info/Debug/Trace 级别可替换的消息 Marshaler消息序列化策略应可配置、可替换参考 Kafka 适配器kafka.Marshaler的设计健壮的Close()发布器与订阅器的 Close 必须幂等在发布器/订阅器被阻塞如等待 Ack时仍能正确关闭在订阅器输出通道无人读取时也能正确关闭完整的 Ack/Nack消费到的消息必须同时支持Ack()与Nack()Nack 后的重投递Nack()后消息应能重新投递这正是 at-least-once 语义的体现参见 docs/content/docs/pub-sub.md跑通通用测试使用 pubsub/tests/test_pubsub.go 中的通用测试套件与压力测试声明实现的Features测试排查技巧见 troubleshooting性能优化文档完备提供 GoDoc、Markdown 文档对应 docs/content/pubsubs/ 目录下的中间件文档格式与 Getting Started 示例对应 _examples/pubsubs/ 目录的格式。以仓库内置的 GoChannel 实现 pubsub/gochannel/pubsub.go 为最小参考它的Config只有四个字段输出通道缓冲OutputChannelBuffer、内存持久化Persistent、发布阻塞直到订阅者 Ack 的BlockPublishUntilSubscriberAck、保留 Context 的PreserveContext构造函数NewGoChannel(config, logger)返回*GoChannel实现了完整的 Publisher/Subscriber。第三方适配器的接入形态与之完全一致——只是把内存 channel 换成外部中间件的连接。接入到路由器后发布侧调用Publish(topic, messages...)消费侧通过Subscribe(ctx, topic)拿到-chan *message.Message配合Ack()/Nack()即可纳入 Watermill 的路由与中间件体系。完成实现后欢迎向官方提交 Pull Request或把自己的库补充进 docs/content/docs/awesome.md 的 Awesome 清单帮助后来的使用者发现它。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表