调用第三方 API 的 Go 服务通常会在内存中缓存 OAuth 访问令牌,并在令牌即将过期时刷新。单实例、低并发时,这段逻辑看起来很简单;并发升高后,几十个请求可能同时发现令牌失效,又同时访问认证端点。结果不是更快拿到新令牌,而是认证服务返回 429 Too Many Requests,业务请求批量失败,甚至出现“较晚完成的旧刷新结果覆盖较新的令牌”。

本文以一个调用支付平台 API 的服务为例,说明如何用 golang.org/x/sync/singleflight 合并同一进程内的令牌刷新,同时保证每个业务请求仍有独立的取消和超时边界。

适用场景

这套方案适用于以下情况:

  • 多个 goroutine 共享同一组 OAuth 客户端凭据和访问令牌;
  • 令牌刷新是昂贵或受限的外部调用,认证端点有严格的频率限制;
  • 业务请求可以共享同一次刷新结果,但不能共享彼此的等待超时;
  • 服务以单进程内存保存令牌,或者每个实例允许独立持有一份令牌;
  • 刷新失败后允许下一批请求再次尝试,而不是长期缓存失败结果。

如果上游要求“每次刷新都会立即撤销上一枚令牌”,而服务有多个副本,就不能只做进程内合并。此时还需要分布式协调,或者由独立的令牌代理统一刷新。后文会单独说明这个边界。

现象描述

生产环境常见的时间线如下:

  1. 访问令牌在 10:00:00 过期;
  2. 10:00:00 附近有 80 个请求同时进入;
  3. 80 个 goroutine 都判断令牌不可用,并同时调用 /oauth/token;
  4. 认证服务只允许每秒刷新 5 次,其余请求返回 429;
  5. 部分刷新请求成功,但完成顺序不同,内存中的令牌被反复覆盖;
  6. 业务日志里同时出现认证失败、限流和请求超时,表面上像是上游整体故障。

典型问题代码如下:

func (m *TokenManager) Get(ctx context.Context) (Token, error) {
	if token := m.current(); token.IsValid(time.Now()) {
		return token, nil
	}

	// 每个并发请求都会进入刷新调用。
	token, err := m.client.Refresh(ctx)
	if err != nil {
		return Token{}, fmt.Errorf("刷新访问令牌: %w", err)
	}
	m.store(token)
	return token, nil
}

只在 store 周围加锁没有用,因为昂贵的网络调用已经并发发生。把整个 Refresh 包在互斥锁中虽然能避免并发刷新,却会让后续请求逐个获得锁并再次刷新;如果锁内不做第二次有效性检查,80 个请求仍可能顺序刷新 80 次。

根因与设计目标

问题的核心是“同一个逻辑操作存在大量同时在途的重复调用”。正确目标不是让所有请求串行执行,而是:

  • 同一时刻最多执行一次令牌刷新;
  • 后来的请求等待并复用这次刷新结果;
  • 等待者取消时只退出自己的等待,不取消其他请求依赖的刷新;
  • 刷新本身有独立且有限的超时,不能无限占用 singleflight 键;
  • 进入刷新函数后再次检查令牌,关闭并发窗口;
  • 成功结果写入内存,失败结果只返回给当前等待者,不长期缓存;
  • 日志不记录 access token、refresh token、客户端密钥或完整响应体。

singleflight.Group 正好提供“相同 key 的在途调用只执行一次”的语义。Do 会让调用者同步等待,DoChan 则返回结果通道,便于同时监听业务请求的 ctx.Done()。

可直接落地的实现

下面的示例使用互斥锁保护令牌快照,用 singleflight.Group 合并刷新。生产项目可以把 TokenClient 替换为真实的 OAuth 客户端。

package token

import (
	"context"
	"errors"
	"fmt"
	"sync"
	"time"

	"golang.org/x/sync/singleflight"
)

// Token 表示认证服务签发的访问令牌及其过期时间。
type Token struct {
	AccessToken string
	ExpiresAt   time.Time
}

// IsValid 判断令牌在预留安全窗口后是否仍然有效。
func (t Token) IsValid(now time.Time, refreshSkew time.Duration) bool {
	return t.AccessToken != "" && now.Add(refreshSkew).Before(t.ExpiresAt)
}

// TokenClient 定义访问认证服务的边界。
type TokenClient interface {
	Refresh(ctx context.Context) (Token, error)
}

// Manager 合并并发刷新,并保存最近一次有效令牌。
type Manager struct {
	client         TokenClient
	refreshTimeout time.Duration
	refreshSkew    time.Duration

	mu    sync.RWMutex
	token Token
	group singleflight.Group
}

// NewManager 创建令牌管理器。
func NewManager(client TokenClient, refreshTimeout, refreshSkew time.Duration) (*Manager, error) {
	if client == nil {
		return nil, errors.New("令牌客户端不能为空")
	}
	if refreshTimeout <= 0 {
		return nil, errors.New("刷新超时必须大于零")
	}
	if refreshSkew < 0 {
		return nil, errors.New("刷新安全窗口不能小于零")
	}

	return &Manager{
		client:         client,
		refreshTimeout: refreshTimeout,
		refreshSkew:    refreshSkew,
	}, nil
}

// Get 返回有效令牌;令牌不可用时合并并发刷新。
func (m *Manager) Get(ctx context.Context) (Token, error) {
	if ctx == nil {
		return Token{}, errors.New("上下文不能为空")
	}
	if token := m.current(); token.IsValid(time.Now(), m.refreshSkew) {
		return token, nil
	}

	resultChannel := m.group.DoChan("oauth_access_token", func() (any, error) {
		// 进入合并函数后再次检查,避免前一个刷新刚完成又重复调用上游。
		if token := m.current(); token.IsValid(time.Now(), m.refreshSkew) {
			return token, nil
		}

		// 刷新不能绑定某个业务请求,否则首个请求取消会让所有等待者失败。
		refreshContext, cancel := context.WithTimeout(context.Background(), m.refreshTimeout)
		defer cancel()

		token, err := m.client.Refresh(refreshContext)
		if err != nil {
			return Token{}, fmt.Errorf("调用认证服务刷新令牌: %w", err)
		}
		if !token.IsValid(time.Now(), m.refreshSkew) {
			return Token{}, errors.New("认证服务返回的令牌缺失或有效期过短")
		}

		m.store(token)
		return token, nil
	})

	select {
	case <-ctx.Done():
		return Token{}, fmt.Errorf("等待令牌刷新: %w", ctx.Err())
	case result := <-resultChannel:
		if result.Err != nil {
			return Token{}, result.Err
		}
		token, ok := result.Val.(Token)
		if !ok {
			return Token{}, errors.New("令牌刷新结果类型异常")
		}
		return token, nil
	}
}

func (m *Manager) current() Token {
	m.mu.RLock()
	defer m.mu.RUnlock()
	return m.token
}

func (m *Manager) store(token Token) {
	m.mu.Lock()
	m.token = token
	m.mu.Unlock()
}

为什么要检查两次

第一次检查是快速路径:令牌有效时完全不进入 singleflight。第二次检查发生在刷新函数内部,用于关闭以下竞态窗口:请求 A 第一次检查发现过期,准备进入合并;请求 B 已经完成刷新并退出;随后 A 成为新一轮调用的执行者。如果没有第二次检查,A 会立刻再次刷新。

为什么使用 DoChan

Do 会一直等到共享刷新完成,业务请求即使已经取消也不能提前返回。DoChan 允许等待者在刷新结果和自己的 ctx.Done() 之间选择。

官方文档还有一个容易误用的细节:DoChan 返回的通道不会关闭,因此这里只接收一次结果,不要对它使用 for range。

为什么刷新不用请求的 Context

假设请求 A 首先进入刷新函数,请求 B 随后复用它。如果刷新直接使用 A 的 ctx,A 的客户端断开连接就会取消共享刷新,B 也只能得到失败结果。

示例使用 context.Background() 创建独立刷新预算,并用 WithTimeout 保证它一定结束。代价是所有等待者都取消后,刷新最多还会继续一个 refreshTimeout。这是有意的取舍:完成刷新可以服务下一批请求,固定超时又能限制资源占用。服务关闭时如果必须立即中断刷新,可以在更高层由专门的令牌代理管理生命周期,但不要把某个随机业务请求当作共享任务的所有者。

调用方如何使用

业务请求仍然用自己的超时等待令牌:

func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
	waitContext, cancel := context.WithTimeout(r.Context(), 800*time.Millisecond)
	defer cancel()

	token, err := h.tokens.Get(waitContext)
	if err != nil {
		h.logger.Error("获取上游访问令牌失败",
			"operation", "get_access_token",
			"error", err,
		)
		http.Error(w, "上游认证暂不可用", http.StatusServiceUnavailable)
		return
	}

	// 令牌只放入请求头,不写入日志、指标标签或错误消息。
	if err := h.callPaymentAPI(r.Context(), token.AccessToken); err != nil {
		h.logger.Error("调用支付接口失败", "error", err)
		http.Error(w, "上游服务异常", http.StatusBadGateway)
	}
}

这里存在两个不同预算:

  • 800ms 是单个业务请求愿意等待令牌的时间;
  • refreshTimeout 是共享刷新允许占用认证服务的最长时间。

两者不能混为一谈。业务请求可以先超时返回,但共享刷新仍可能成功,随后到来的请求会直接命中内存令牌。

并发测试:确认只刷新一次

下面的测试让 50 个 goroutine 同时请求令牌,并统计认证客户端被调用的次数。

package token

import (
	"context"
	"sync"
	"sync/atomic"
	"testing"
	"time"
)

type fakeClient struct {
	calls atomic.Int64
	gate  <-chan struct{}
}

func (c *fakeClient) Refresh(ctx context.Context) (Token, error) {
	c.calls.Add(1)
	select {
	case <-ctx.Done():
		return Token{}, ctx.Err()
	case <-c.gate:
		return Token{
			AccessToken: "仅用于测试的令牌",
			ExpiresAt:   time.Now().Add(time.Hour),
		}, nil
	}
}

func TestManagerMergesConcurrentRefresh(t *testing.T) {
	gate := make(chan struct{})
	client := &fakeClient{gate: gate}
	manager, err := NewManager(client, time.Second, time.Minute)
	if err != nil {
		t.Fatalf("创建令牌管理器失败: %v", err)
	}

	const workers = 50
	start := make(chan struct{})
	errorsChannel := make(chan error, workers)
	var waitGroup sync.WaitGroup

	for range workers {
		waitGroup.Add(1)
		go func() {
			defer waitGroup.Done()
			<-start
			_, getErr := manager.Get(context.Background())
			errorsChannel <- getErr
		}()
	}

	close(start)
	deadline := time.Now().Add(time.Second)
	for client.calls.Load() == 0 && time.Now().Before(deadline) {
		time.Sleep(time.Millisecond)
	}
	close(gate)
	waitGroup.Wait()
	close(errorsChannel)

	for getErr := range errorsChannel {
		if getErr != nil {
			t.Fatalf("获取令牌失败: %v", getErr)
		}
	}
	if calls := client.calls.Load(); calls != 1 {
		t.Fatalf("期望只刷新一次,实际刷新 %d 次", calls)
	}
}

执行:

go test -race -count=20 ./...
  • -race 检查令牌读写和测试协程之间是否存在数据竞争;
  • -count=20 增加不同调度顺序被覆盖的概率;
  • 测试应断言上游调用次数,而不只断言 50 个请求都成功。

还应补充三个失败路径:认证服务返回错误时所有等待者收到同一错误;某个等待者超时时其他等待者仍能成功;认证服务返回空令牌或过短有效期时不得写入内存。

可观测性怎么做

建议记录和监控以下字段:

  • token_refresh_total{result="success|error"}:真实刷新次数;
  • token_refresh_duration_seconds:认证端点耗时;
  • token_wait_timeout_total:业务请求等待共享结果超时次数;
  • token_expires_in_seconds:当前令牌剩余有效期;
  • shared:singleflight.Result.Shared,用于确认请求是否发生合并。

固定日志文案可以写“刷新访问令牌失败”,动态字段保留 provider、status_code、duration_ms、trace_id 和错误链。禁止记录令牌、客户端密钥、Authorization 头、Cookie 或认证服务的完整原始响应。

监控时不要把令牌值、用户 ID 等高基数数据放进指标标签。需要关联问题时使用受控的 trace_id,并确保链路系统同样不会采集敏感请求头。

常见误区

把 singleflight 当缓存

singleflight 只合并同时在途的调用。刷新结束后,该 key 不会保存结果;下一次调用是否命中,取决于令牌是否已经写入自己的存储。因此快速路径和刷新函数内的二次检查都不能省略。

为每个请求创建一个 Group

Group 必须被同一批并发请求共享。若在 Get 内部创建局部 Group,每个请求仍然各自刷新,等于没有合并。

超时后立即调用 Forget

等待者超时不代表共享刷新已经失效。此时调用 Forget 会允许新请求启动第二个并发刷新,重新制造雪崩。只有业务明确判断当前在途操作必须被替换,并能承受并发执行时,才考虑 Forget。

所有租户共用一个固定 key

如果系统按租户、区域或 OAuth 客户端分别持有令牌,key 必须包含这些隔离维度。否则租户 A 可能等待租户 B 的刷新,甚至拿到错误的凭据结果。拼接 key 时要使用无歧义编码,例如长度前缀或结构化摘要,不要直接用可能碰撞的字符串相加。

认为进程内合并等于集群级合并

10 个 Pod 各自有一个 Group,令牌过期时仍可能产生 10 次刷新。若认证服务允许多个有效令牌,这通常可接受;若刷新会撤销旧令牌,则应使用单独的令牌服务、数据库行锁或带所有权令牌的分布式锁。分布式锁还必须处理租约过期、进程暂停和旧持有者覆盖新结果,不能只写一个 SET NX 就认为完成了协调。

预防与上线检查

上线前至少确认以下事项:

  1. 过期判断包含合理的安全窗口,避免令牌在请求飞行途中失效;
  2. 快速路径与刷新函数内部都检查令牌有效性;
  3. Group 的生命周期覆盖所有需要合并的并发请求;
  4. 业务等待超时与共享刷新超时彼此独立;
  5. 刷新有明确的网络超时、退出码或状态码检查和错误包装;
  6. 上游返回的令牌非空,过期时间满足最小有效期;
  7. 失败结果不会覆盖上一枚仍可使用的令牌;
  8. 单元测试覆盖并发合并、等待者取消、刷新失败和非法响应;
  9. CI 执行 go test -race ./...;
  10. 日志、指标和追踪系统均不采集任何认证秘密。

总结

OAuth 令牌刷新雪崩不是简单的“加一把锁”问题,而是共享操作的所有权、等待者取消和结果发布边界没有设计清楚。singleflight 可以把同一进程内的重复在途刷新合并为一次,但它不负责缓存,也不会自动处理业务超时或跨实例协调。

可靠实现应同时具备四点:有效令牌走快速路径;共享刷新内部做二次检查;每个调用方独立等待和取消;刷新任务使用有限且独立的执行预算。再配合并发调用次数断言、竞态检测和敏感日志审计,才能让令牌过期从一次流量尖峰事故,变成可预测、可验证的正常状态切换。