From 4fe9f83b89afc2cb10e41b15bc3d4ded73f28b33 Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 15 Aug 2026 17:49:31 +0800 Subject: [PATCH] =?UTF-8?q?chore(=E8=AE=A2=E9=98=85):=20=E9=87=8D=E8=BF=9E?= =?UTF-8?q?=E9=94=99=E8=AF=AF=E6=97=A5=E5=BF=97=E8=A1=A5=E5=85=85=20Broker?= =?UTF-8?q?=20=E5=9C=B0=E5=9D=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Motivation: 多 Broker 部署场景下,订阅者重连失败时仅记录错误名称,难以快速定位是哪一台 Broker 连接异常。补充 Broker 地址并规范化错误字段,提升故障排查效率。 Changes: * 重连失败日志新增 broker host 字段 * 新增 brokerHost 函数,从连接 URL 解析并输出主机地址 * 错误字段改为记录错误信息字符串,避免直接序列化 error 对象 --- internal/subscriber/subscriber.go | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/internal/subscriber/subscriber.go b/internal/subscriber/subscriber.go index 34bb797..fcb2d23 100644 --- a/internal/subscriber/subscriber.go +++ b/internal/subscriber/subscriber.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "log/slog" + "net/url" "time" "aiaa-notification-service/internal/config" @@ -42,7 +43,10 @@ func (s *Subscriber) Run(ctx context.Context) error { } if err := s.consumeOnce(ctx); err != nil { - slog.Error("subscriber error, reconnecting", "name", s.cfg.Name, "error", err) + slog.Error("subscriber error, reconnecting", + "name", s.cfg.Name, + "host", brokerHost(s.cfg.URL), + "error", err.Error()) select { case <-ctx.Done(): return ctx.Err() @@ -159,6 +163,17 @@ func (s *Subscriber) republish(ch *amqp.Channel, d amqp.Delivery, queue string) slog.Info("message requeued", "name", s.cfg.Name, "retry", headers[retryHeader], "max", s.cfg.MaxRetry) } +func brokerHost(raw string) string { + u, err := url.Parse(raw) + if err != nil || u.Hostname() == "" { + return "" + } + if u.Port() != "" { + return u.Hostname() + ":" + u.Port() + } + return u.Hostname() +} + func copyAMQPHeaders(headers amqp.Table) amqp.Table { out := amqp.Table{} for k, v := range headers {