chore(订阅): 重连错误日志补充 Broker 地址
Motivation: 多 Broker 部署场景下,订阅者重连失败时仅记录错误名称,难以快速定位是哪一台 Broker 连接异常。补充 Broker 地址并规范化错误字段,提升故障排查效率。 Changes: * 重连失败日志新增 broker host 字段 * 新增 brokerHost 函数,从连接 URL 解析并输出主机地址 * 错误字段改为记录错误信息字符串,避免直接序列化 error 对象
This commit is contained in:
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user