/
githubmirror
/
ydb-go-sdk
Обзор
Документация
Войти
/
githubmirror
/
ydb-go-sdk
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
internal/query/session.go
288 строк
7 KB
Aleksey Myasnikov
Deprecate WithConcurrentResultSets and always enable it in Query (#2216)
29 июн 2026, 15:49
Не верифицирован
29 июн 2026, 15:49
3810a99
Код
Авторство
О чём код?
package query import ( "context" "time" "github.com/ydb-platform/ydb-go-genproto/Ydb_Query_V1" "github.com/ydb-platform/ydb-go-genproto/protos/Ydb" "github.com/ydb-platform/ydb-go-sdk/v3/internal/query/arrow" "github.com/ydb-platform/ydb-go-sdk/v3/internal/query/gtrace" "github.com/ydb-platform/ydb-go-sdk/v3/internal/query/options" "github.com/ydb-platform/ydb-go-sdk/v3/internal/query/result" "github.com/ydb-platform/ydb-go-sdk/v3/internal/stack" baseTx "github.com/ydb-platform/ydb-go-sdk/v3/internal/tx" "github.com/ydb-platform/ydb-go-sdk/v3/internal/xcontext" "github.com/ydb-platform/ydb-go-sdk/v3/internal/xerrors" "github.com/ydb-platform/ydb-go-sdk/v3/query" "github.com/ydb-platform/ydb-go-sdk/v3/trace" ) var _ query.Session = (*Session)(nil) type ( Session struct { Core client Ydb_Query_V1.QueryServiceClient trace *trace.Query lazyTx bool streamResultCloseTimeout time.Duration } ) func (s *Session) QueryResultSet( ctx context.Context, q string, opts ...options.Execute, ) (rs result.ClosableResultSet, finalErr error) { settings := options.ExecuteSettings(opts...) if err := validateTxControl(settings); err != nil { return nil, err } onDone := gtrace.QueryOnSessionQueryResultSet(s.trace, &ctx, stack.FunctionID("github.com/ydb-platform/ydb-go-sdk/v3/internal/query.(*Session).QueryResultSet"), s, q) defer func() { onDone(finalErr) }() r, err := s.execute(ctx, q, settings, options.ResultSetsTypeOrdered, withStreamResultTrace(s.trace), withIssuesHandler(settings.IssuesOpts()), ) if err != nil { return nil, xerrors.WithStackTrace(err) } rs, err = readResultSet(ctx, r) if err != nil { return nil, xerrors.WithStackTrace(err) } return rs, nil } func (s *Session) queryRow( ctx context.Context, q string, settings executeSettings, resultOpts ...resultOption, ) (row query.Row, finalErr error) { r, err := s.execute(ctx, q, settings, options.ResultSetsTypeOrdered, resultOpts...) if err != nil { return nil, xerrors.WithStackTrace(err) } defer func() { _ = r.Close(ctx) }() row, err = readRow(ctx, r) if err != nil { return nil, xerrors.WithStackTrace(err) } return row, nil } func (s *Session) QueryRow(ctx context.Context, q string, opts ...options.Execute) (_ query.Row, finalErr error) { settings := options.ExecuteSettings(opts...) if err := validateTxControl(settings); err != nil { return nil, err } onDone := gtrace.QueryOnSessionQueryRow(s.trace, &ctx, stack.FunctionID("github.com/ydb-platform/ydb-go-sdk/v3/internal/query.(*Session).QueryRow"), s, q) defer func() { onDone(finalErr) }() row, err := s.queryRow(ctx, q, settings, withStreamResultTrace(s.trace), withIssuesHandler(settings.IssuesOpts()), ) if err != nil { return nil, xerrors.WithStackTrace(err) } return row, nil } func createSession( ctx context.Context, client Ydb_Query_V1.QueryServiceClient, opts ...Option, ) (*Session, error) { core, err := Open(ctx, client, opts...) if err != nil { return nil, xerrors.WithStackTrace(err) } return &Session{ Core: core, trace: core.Trace, client: core.Client, streamResultCloseTimeout: core.deleteTimeout, }, nil } func (s *Session) Begin( ctx context.Context, txSettings query.TransactionSettings, ) ( tx query.Transaction, finalErr error, ) { onDone := gtrace.QueryOnSessionBegin(s.trace, &ctx, stack.FunctionID("github.com/ydb-platform/ydb-go-sdk/v3/internal/query.(*Session).Begin"), s) defer func() { if finalErr != nil { applyStatusByError(s, finalErr) onDone(finalErr, nil) } else { onDone(nil, tx) } }() if lazyTx := baseTx.LazyTxFromContext(ctx, s.lazyTx); lazyTx { if !s.IsAlive() { return nil, xerrors.WithStackTrace(xerrors.Operation( xerrors.WithStatusCode(Ydb.StatusIds_BAD_SESSION), )) } return &Transaction{ s: s, txSettings: txSettings, }, nil } txID, err := begin(ctx, s, txSettings) if err != nil { return nil, xerrors.WithStackTrace(err) } return &Transaction{ LazyID: baseTx.ID(txID), txSettings: txSettings, s: s, }, nil } func (s *Session) execute(ctx context.Context, q string, settings executeSettings, concurrentResultSets options.ResultSetsType, opts ...resultOption, ) (_ *streamResult, finalErr error) { ctx, cancel := context.WithCancel(ctx) defer func() { if finalErr != nil { cancel() applyStatusByError(s, finalErr) } }() r, err := execute(ctx, s.ID(), s.client, q, settings, concurrentResultSets, append(opts, withStreamResultCloseTimeout(s.streamResultCloseTimeout), withStreamResultOnClose(cancel), )...) if err != nil { return nil, xerrors.WithStackTrace(err) } return r, nil } func (s *Session) Exec(ctx context.Context, q string, opts ...options.Execute) (finalErr error) { settings := options.ExecuteSettings(opts...) if err := validateTxControl(settings); err != nil { return err } onDone := gtrace.QueryOnSessionExec(s.trace, &ctx, stack.FunctionID("github.com/ydb-platform/ydb-go-sdk/v3/internal/query.(*Session).Exec"), s, q, settings.Label(), ) defer func() { onDone(finalErr) }() r, err := s.execute(ctx, q, settings, options.ResultSetsTypeOrdered, withStreamResultTrace(s.trace), withIssuesHandler(settings.IssuesOpts()), ) if err != nil { return xerrors.WithStackTrace(err) } defer func() { _ = r.Close(ctx) }() err = readAll(ctx, r) if err != nil { return xerrors.WithStackTrace(err) } return nil } func (s *Session) Query(ctx context.Context, q string, opts ...options.Execute) (_ query.Result, finalErr error) { settings := options.ExecuteSettings(opts...) if err := validateTxControl(settings); err != nil { return nil, err } onDone := gtrace.QueryOnSessionQuery(s.trace, &ctx, stack.FunctionID("github.com/ydb-platform/ydb-go-sdk/v3/internal/query.(*Session).Query"), s, q, settings.Label(), ) defer func() { onDone(finalErr) }() r, err := s.execute(ctx, q, settings, options.ResultSetsTypeOrdered, withStreamResultTrace(s.trace), withIssuesHandler(settings.IssuesOpts()), ) if err != nil { return nil, xerrors.WithStackTrace(err) } return r, nil } // QueryArrow like [*Session.Query] but returns results in [Apache Arrow] format. // Each part of the result implements io.Reader and contains the data in Arrow IPC format. // // Experimental: https://github.com/ydb-platform/ydb-go-sdk/blob/master/VERSIONING.md#experimental // // [Apache Arrow]: https://arrow.apache.org/ func (s *Session) QueryArrow(ctx context.Context, q string, opts ...options.Execute) (_ arrow.Result, finalErr error) { settings := options.ExecuteSettings(opts...) defer func() { if finalErr != nil { applyStatusByError(s, finalErr) } }() request, callOptions, err := executeQueryRequest(s.ID(), q, settings, options.ResultSetsTypeOrdered) if err != nil { return nil, xerrors.WithStackTrace(err) } request.ResultSetFormat = Ydb.ResultSet_FORMAT_ARROW executeCtx, executeCancel := context.WithCancel(xcontext.ValueOnly(ctx)) defer func() { if finalErr != nil { executeCancel() } }() stream, err := s.client.ExecuteQuery(executeCtx, request, callOptions...) if err != nil { return nil, xerrors.WithStackTrace(err) } return &arrowResult{stream: stream, close: executeCancel}, nil }