/
githubmirror
/
ydb-go-sdk
Обзор
Документация
Войти
/
githubmirror
/
ydb-go-sdk
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
internal/query/execute_query_stream_prefetch.go
88 строк
2 KB
Aleksey Myasnikov
perf: prefetch N parts of QueryClient response stream for decrease query latency (#2150)
12 май 2026, 18:58
Не верифицирован
12 май 2026, 18:58
3455d18
Код
Авторство
О чём код?
package query import ( "fmt" "io" "github.com/ydb-platform/ydb-go-genproto/Ydb_Query_V1" "github.com/ydb-platform/ydb-go-genproto/protos/Ydb_Query" "google.golang.org/protobuf/proto" "github.com/ydb-platform/ydb-go-sdk/v3/internal/xerrors" ) // executeQueryPartRecv is one logical item from ExecuteQuery stream (typed Recv). type executeQueryPartRecv struct { part *Ydb_Query.ExecuteQueryResponsePart err error } type asyncPrefetchExecuteQueryStream struct { Ydb_Query_V1.QueryService_ExecuteQueryClient ch chan executeQueryPartRecv } func wrapExecuteQueryStreamWithAsyncPrefetch( stream Ydb_Query_V1.QueryService_ExecuteQueryClient, prefetch int, ) Ydb_Query_V1.QueryService_ExecuteQueryClient { if prefetch <= 0 { return stream } s := &asyncPrefetchExecuteQueryStream{ QueryService_ExecuteQueryClient: stream, ch: make(chan executeQueryPartRecv, prefetch), } go s.pump() return s } func (p *asyncPrefetchExecuteQueryStream) pump() { defer close(p.ch) ctx := p.QueryService_ExecuteQueryClient.Context() for { part, err := p.QueryService_ExecuteQueryClient.Recv() item := executeQueryPartRecv{part: part, err: err} select { case p.ch <- item: case <-ctx.Done(): return } if err != nil { return } } } func (p *asyncPrefetchExecuteQueryStream) Recv() (*Ydb_Query.ExecuteQueryResponsePart, error) { item, ok := <-p.ch if !ok { return nil, io.EOF } return item.part, item.err } func (p *asyncPrefetchExecuteQueryStream) RecvMsg(m any) error { part, err := p.Recv() if err != nil { return err } dst, ok := m.(*Ydb_Query.ExecuteQueryResponsePart) if !ok { return xerrors.WithStackTrace(fmt.Errorf( "%T is not '*Ydb_Query.ExecuteQueryResponsePart'", m, )) } proto.Reset(dst) if part != nil { proto.Merge(dst, part) } return nil }