You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

测试分页查询取消功能时的时序精度问题处理

问题描述

我实现了一个用于执行异步分页SQL查询的PagedQuery对象,代码如下:

type PagedQuery[T any] struct {
    Results   chan []*T
    Errors    chan error
    Done      chan error
    Quit      chan error
    client    *sql.DB
}

func NewPagedQuery[T any](client *sql.DB) *PagedQuery[T] {
    return &PagedQuery[T]{
        Results:   make(chan []*T, 1),
        Errors:    make(chan error, 1),
        Done:      make(chan error, 1),
        Quit:      make(chan error, 1),
        client:    client,
    }
}

func (paged *PagedQuery[T]) requestAsync(ctx context.Context, queries ...*Query) {

    conn, err := client.Conn(ctx)
    if err != nil {
        paged.Errors <- err
        return
    }

    defer func() {
        conn.Close()
        paged.Done <- nil
    }()

    for i, query := range queries {
        select {
        case <-ctx.Done():
            return
        case <-paged.Quit:
            return
        default:
        }

        rows, err := conn.QueryContext(ctx, query.String, query.Arguments...)
        if err != nil {
            paged.Errors <- err
            return
        }

        data, err := sql.ReadRows[T](rows)
        if err != nil {
            paged.Errors <- err
            return
        }

        paged.Results <- data
    }
}

我正测试该对象的查询取消功能,测试代码如下:

svc, mock := createServiceMock("TEST_DATABASE", "TEST_SCHEMA")

mock.ExpectQuery(regexp.QuoteMeta("TEST QUERY")).
    WithArgs(...).
    WillReturnRows(mock.NewRows([]string{"t", "v", "o", "c", "h", "l", "vw", "n"}))

ctx, cancel := context.WithCancel(context.Background())
go svc.requestAsync(ctx, query1, query2, query3, query4)

time.Sleep(50 * time.Millisecond)
cancel()

results := make([]data, 0)
loop:
for {
    select {
    case <-query.Done:
        break loop
    case err := <-query.Errors:
        Expect(err).ShouldNot(HaveOccurred())
    case r := <-query.Results:
        results = append(results, r...)
    }
}

Expect(results).Should(BeEmpty())
Expect(mock.ExpectationsWereMet()).ShouldNot(HaveOccurred()) // 此处偶尔失败

问题在于测试偶尔会在标注行失败:调用cancel()时,代码不一定处于检查ctx.Done()或Quit通道的select分支,可能处于循环内执行查询或发送结果的阶段。我原本认为执行会阻塞到接收Results通道数据后才继续,但实际并非如此。另外我使用sqlmock进行SQL测试,它不支持SQL查询的模糊校验。请问该失败的原因是什么,如何修复?


失败原因分析
  • 取消时机不可靠:依赖time.Sleep(50ms)同步属于完全不可控的方式,goroutine可能在sleep期间已经执行了1次甚至多次查询,导致sqlmock的期望(比如仅预期1次查询)被打破。
  • 带缓冲的Results通道:Results通道缓冲大小为1,发送数据时不会阻塞,goroutine会直接进入下一轮循环,不会等待接收方处理结果,因此即使触发cancel,goroutine可能已经完成多轮查询。
  • 取消信号检查不完整:仅在循环开头的select分支检查取消信号,conn.QueryContext、sql.ReadRows、paged.Results <- data这些阶段没有实时响应ctx.Done,导致取消后这些操作仍会继续执行。
  • sqlmock期望精确匹配:sqlmock的查询期望是精确计数的,如果goroutine执行了超出预期次数的查询,ExpectationsWereMet()必然失败。

修复方案

1. 替换不可靠的time.Sleep,改用同步机制控制取消时机

用通道等待goroutine进入执行流程后再触发取消,避免因sleep时长不合适导致的时序问题:

ctx, cancel := context.WithCancel(context.Background())
startChan := make(chan struct{})
go func() {
    defer close(startChan)
    svc.requestAsync(ctx, query1, query2, query3, query4)
}()
<-startChan // 等待goroutine启动后再取消
cancel()

2. 修改PagedQuery,在关键阶段添加取消检查

调整requestAsync方法,在所有耗时/阻塞操作前后添加取消检查,同时将Results改为无缓冲通道,确保发送结果时能感知取消信号:

func NewPagedQuery[T any](client *sql.DB) *PagedQuery[T] {
    return &PagedQuery[T]{
        Results:   make(chan []*T), // 改为无缓冲通道
        Errors:    make(chan error, 1),
        Done:      make(chan struct{}), // 改为信号通道,关闭表示完成
        Quit:      make(chan struct{}),
        client:    client,
    }
}

func (paged *PagedQuery[T]) requestAsync(ctx context.Context, queries ...*Query) {
    // 修复原代码的错误:引用paged.client而非全局client
    conn, err := paged.client.Conn(ctx)
    if err != nil {
        paged.Errors <- err
        close(paged.Done)
        return
    }

    defer func() {
        conn.Close()
        close(paged.Done) // 用关闭通道代替发送nil,避免重复发送
    }()

    for _, query := range queries {
        // 循环开头检查取消
        select {
        case <-ctx.Done():
            return
        case <-paged.Quit:
            return
        default:
        }

        rows, err := conn.QueryContext(ctx, query.String, query.Arguments...)
        if err != nil {
            paged.Errors <- err
            return
        }

        // 查询完成后立即检查取消,避免继续处理结果
        select {
        case <-ctx.Done():
            rows.Close()
            return
        case <-paged.Quit:
            rows.Close()
            return
        default:
        }

        data, err := sql.ReadRows[T](rows)
        rows.Close() // 及时关闭rows释放资源
        if err != nil {
            paged.Errors <- err
            return
        }

        // 发送结果时检查取消,无缓冲通道会阻塞直到接收方处理
        select {
        case <-ctx.Done():
            return
        case <-paged.Quit:
            return
        case paged.Results <- data:
            // 发送成功后再次检查取消
            select {
            case <-ctx.Done():
                return
            case <-paged.Quit:
                return
            default:
            }
        }
    }
}

3. 优化sqlmock的期望设置

如果测试的是取消功能,允许查询被执行0次或1次,使用Maybe()方法放宽期望:

mock.ExpectQuery(regexp.QuoteMeta("TEST QUERY")).
    WithArgs(...).
    WillReturnRows(mock.NewRows([]string{"t", "v", "o", "c", "h", "l", "vw", "n"})).
    Maybe() // 允许该期望被执行0次或多次

4. 调整测试循环,正确处理通道关闭

原测试中Done通道的处理逻辑不完善,改为检查通道是否关闭:

results := make([]data, 0)
loop:
for {
    select {
    case _, ok := <-svc.Done:
        if !ok {
            break loop
        }
    case err := <-svc.Errors:
        Expect(err).ShouldNot(HaveOccurred())
    case r := <-svc.Results:
        results = append(results, r...)
    }
}

内容的提问来源于stack exchange,提问作者Woody1193

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 23:50:34