spaniter

package module
v0.3.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jun 27, 2026 License: MIT Imports: 7 Imported by: 0

README

spaniter

Go Reference

spaniter adapts Cloud Spanner row streams to Go standard iterators.

The module is deliberately lower-level than github.com/apstndb/spanvalue: it does not format values, choose output formats, or own export policy. Its job is to turn *spanner.RowIterator into iter.Seq2[*spanner.Row, error] while preserving the Spanner iterator lifecycle.

API

  • RowIteratorSeq: adapts *spanner.RowIterator to iter.Seq2[*spanner.Row, error].
  • DrainRowIterator: consumes *spanner.RowIterator without yielding rows and returns metadata/stats.
  • WithResult: captures metadata, stats, and rows read in a RowIteratorResult.
  • WithOnMetadata: invokes a callback once when ResultSetMetadata becomes available.
  • WithOnStats: invokes a callback after completion with query plan, query stats, and DML row count.
  • WithStatsEncoding: selects row-count protobuf encoding for drained results.
  • RowIteratorResult.StatsProto: converts captured stats to *sppb.ResultSetStats using the encoding from drain options. Returns nil, nil when there are no top-level ResultSetStats fields to encode. Use only on values returned or captured by spaniter; manual struct literals omit unexported lifecycle state.
  • RowIteratorResult.ResultSet: builds *sppb.ResultSet from materialized rows and captured lifecycle data.
  • PullRowIteratorSeq: iter.Pull2 adapter that normalizes terminal errors.
  • WithDrainOnEarlyStop: optionally drains remaining rows after an early consumer stop so stats can be populated.
  • Rows: adapt already-built rows for tests and virtual result sets.
  • SliceToRowSeq: adapt existing []*spanner.Row fixtures for downstream tests and fakes.

RowIterator lifecycle

Cloud Spanner result metadata becomes available after the first Next call, unless that call returns an error other than iterator.Done. Query plan and query stats become available after Next returns iterator.Done only for an iterator created with QueryWithStats; DML row count becomes available after iterator.Done when the RowIterator represents a DML statement.

RowIteratorSeq keeps those rules explicit while leaving metadata capture optional:

rows := spaniter.RowIteratorSeq(rowIter)

for row, err := range rows {
	if err != nil {
		return err
	}
	_ = row
}

Once the returned sequence is invoked, RowIteratorSeq owns the RowIterator for that run and calls Stop before returning. The sequence is single-use and must not be consumed concurrently. For short call sites, code can pass a freshly created iterator directly instead of binding it only to defer Stop:

stmt := spanner.Statement{SQL: sql}

for row, err := range spaniter.RowIteratorSeq(txn.Query(ctx, stmt)) {
	if err != nil {
		return err
	}
	_ = row
}

The same ownership transfer applies when passing the sequence to a function that immediately consumes iter.Seq2[*spanner.Row, error]. Merely accepting, storing, or forwarding the sequence does not invoke it or stop the underlying RowIterator. Many consumers do not need ResultSetMetadata; they can process rows directly. Use the inline form only when the called function consumes the sequence synchronously before returning; otherwise bind the RowIterator and keep responsibility for Stop.

func consumeRows(rows iter.Seq2[*spanner.Row, error]) error {
	for row, err := range rows {
		if err != nil {
			return err
		}
		if err := processRow(row); err != nil {
			return err
		}
	}
	return nil
}

stmt := spanner.Statement{SQL: sql}

if err := consumeRows(spaniter.RowIteratorSeq(txn.Query(ctx, stmt))); err != nil {
	return err
}

Each yielded pair is either a non-nil row with a nil error, or a nil row with a non-nil terminal error. After yielding a non-nil error, the sequence stops; code should stop processing on the first error.

If the returned lazy sequence is never invoked, callers must still retain responsibility for stopping the original RowIterator.

Use WithResult when code needs result-set metadata, stats, or rows read outside the row loop:

var result spaniter.RowIteratorResult
for row, err := range spaniter.RowIteratorSeq(rowIter, spaniter.WithResult(&result)) {
	if err != nil {
		return err
	}
	_ = row
}
_ = result.Metadata
_ = result.Stats
_ = result.RowsRead

Use WithStatsEncoding when draining if the caller knows the result came from executed standard DML and needs row_count_exact: 0 preserved. The default encoding omits zero row counts because RowIterator exposes row count as a plain int64 and cannot distinguish absent row count from exact zero. Do not use StatsEncodingDMLExact for PLAN mode or read-only queries.

result, err := spaniter.DrainRowIterator(rowIter,
	spaniter.WithStatsEncoding(spaniter.StatsEncodingDMLExact),
)
if err != nil {
	return err
}
stats, err := result.StatsProto()
rs, err := result.ResultSet(rows)

StatsProto and ResultSet depend on unexported state set during draining. Use the value returned from DrainRowIterator or a WithResult capture target; do not build RowIteratorResult manually for protobuf conversion. Callers that only need the Stats struct can keep using WithOnStats without protobuf conversion.

StatsProto returns nil, nil when there are no top-level ResultSetStats fields to encode. Partitioned DML counts come from Client.PartitionedUpdate, not RowIterator.

Use WithOnMetadata or WithOnStats when code needs hook-style callbacks instead of a captured result value. For fan-in adapters or streaming sinks that must publish row-type metadata before the first downstream row, prefer WithOnMetadata: WithResult is updated at the same lifecycle points, but an empty result set never enters the consumer's loop body, so loop-local polling cannot observe metadata until after the sequence is exhausted.

If application code needs metadata or stats but does not want row values, DrainRowIterator consumes the stream internally and returns only lifecycle results:

result, err := spaniter.DrainRowIterator(rowIter)
if err != nil {
	return err
}
_ = result.Metadata
_ = result.Stats
_ = result.RowsRead

This still reads the Spanner result stream internally because the Go client populates stats only after Next returns iterator.Done. To avoid reading data rows at the query level, execute a statement that returns no data rows. If draining returns an error, the returned result can still contain partial metadata and RowsRead; stats are populated only after a successful full drain.

If the consumer may stop early but the caller still needs stats, enable WithDrainOnEarlyStop. This drains after any early stop, including a range-loop break or an adapter that stops pulling because downstream work failed, so keep it disabled unless that extra read work is acceptable. It has no effect on DrainRowIterator, which already drains. When this option is enabled, RowsRead includes rows read during the post-stop drain, including rows not yielded to the consumer. If the post-stop drain fails, its error cannot be yielded because the consumer has already stopped. In that case, WithOnStats is not called and RowIteratorResult.Stats remains zero.

rows := spaniter.RowIteratorSeq(rowIter,
	spaniter.WithDrainOnEarlyStop(),
	spaniter.WithOnStats(func(got spaniter.Stats) {
		stats = got
	}),
)

PullRowIteratorSeq

Consumers using iter.Pull2 should prefer PullRowIteratorSeq. It normalizes terminal errors: when the sequence yields (nil, err), pull returns (nil, err, false) instead of (nil, err, true). Check err before ok.

pull, stop := spaniter.PullRowIteratorSeq(rowIter)
defer stop()

for {
	row, err, ok := pull()
	if err != nil {
		return err
	}
	if !ok {
		break
	}
	_ = row
}

stop releases the RowIterator even when called before the first pull. After the first pull, RowIteratorSeq owns the iterator until the sequence ends.

Example: spanvalue/writer

Applications that already expose rows as iter.Seq2 can compose spaniter with spanvalue/writer; this is caller-side composition, not a dependency of spaniter. When spanvalue/writer is the only consumer of a raw *spanner.RowIterator, its WriteRowIterator helper is the simpler direct path. The composition below is useful when the surrounding pipeline is already expressed as iter.Seq2.

RunRowSeqDeferredMetadata evaluates its metadata function after the first pull, or after an empty sequence ends, so WithResult has already recorded the metadata. RowIteratorHooksFromWriter then registers the row type and flushes after a successful run; for a delimited writer with headers, this permits header-only output for an empty result set. This call synchronously pulls the sequence before returning, so it can take lifecycle ownership of the inline RowIteratorSeq argument.

stmt := spanner.Statement{SQL: sql}

var spannerResult spaniter.RowIteratorResult
if _, err := writer.RunRowSeqDeferredMetadata(
	func() *sppb.ResultSetMetadata { return spannerResult.Metadata },
	spaniter.RowIteratorSeq(
		txn.Query(ctx, stmt),
		spaniter.WithResult(&spannerResult),
	),
	writer.RowIteratorHooksFromWriter(w),
); err != nil {
	return err
}

When metadata is already available and code needs hook-level control, RunRowSeq is the direct consumer form:

if _, err := writer.RunRowSeq(
	metadata,
	spaniter.Rows(row1, row2),
	writer.RowIteratorHooksFromWriter(w),
); err != nil {
	return err
}

For ordinary writer use, WriteRowSeq wraps the same hook setup:

if _, err := writer.WriteRowSeq(
	metadata,
	spaniter.Rows(row1, row2),
	w,
); err != nil {
	return err
}

This keeps spaniter independent of formatting packages while letting callers reuse the standard iterator stream directly. In this sequence-based composition, the writer.RowIteratorResult returned by RunRowSeq, RunRowSeqDeferredMetadata, or WriteRowSeq has zero Stats; use spaniter.WithResult or WithOnStats for Spanner-specific execution data. The writer result's RowsRead counts successful writes, whereas spaniter.RowIteratorResult.RowsRead counts rows consumed from the Spanner iterator.

Development

make check

Documentation

Overview

Package spaniter adapts Cloud Spanner row streams to Go standard iterators.

The package is intentionally lower-level than formatters and writers: it owns only iterator lifecycle concerns such as RowIterator.Stop, result metadata, and post-drain query stats. Formatting, headers, and export policy stay in callers.

Index

Constants

This section is empty.

Variables

View Source
var ErrNilRow = errors.New("nil row")

ErrNilRow reports that an adapted source produced a nil row with a nil error.

View Source
var ErrNilRowIterator = errors.New("nil row iterator")

ErrNilRowIterator reports that RowIteratorSeq was given a nil iterator.

Because RowIteratorSeq returns an iter.Seq2, the error is yielded when the sequence is consumed rather than returned by the constructor.

Functions

func PullRowIteratorSeq added in v0.3.0

func PullRowIteratorSeq(rowIter *spanner.RowIterator, opts ...Option) (pull func() (*spanner.Row, error, bool), stop func())

PullRowIteratorSeq adapts RowIteratorSeq for consumers using iter.Pull2.

The returned pull function normalizes terminal errors from the sequence: when RowIteratorSeq yields (nil, err), pull returns (nil, err, false) instead of (nil, err, true). This matches the usual "check err before ok" pattern and avoids treating iterator failures as EOF.

The returned stop function releases the RowIterator. If stop runs before the first pull, stop calls RowIterator.Stop directly because iter.Pull2 has not started the sequence yet. After the first pull, stop signals the sequence and RowIteratorSeq owns Stop when the sequence goroutine exits.

func RowIteratorSeq

func RowIteratorSeq(rowIter *spanner.RowIterator, opts ...Option) iter.Seq2[*spanner.Row, error]

RowIteratorSeq adapts a cloud.google.com/go/spanner.RowIterator to a Go standard iterator.

The returned sequence owns rowIter: once iteration starts it always calls *cloud.google.com/go/spanner.RowIterator.Stop before returning. Metadata and stats are exposed through WithOnMetadata and WithOnStats hooks instead of requiring callers to keep reading fields from the stopped RowIterator. The sequence is single-use and not safe for concurrent consumption; construct a new RowIterator for another pass.

If the returned sequence is never invoked, RowIteratorSeq cannot call Stop. After constructing a sequence, callers must either consume it, pass it to code that will consume or stop it, or retain responsibility for stopping the original RowIterator.

Each yielded pair is either a non-nil row with a nil error, or a nil row with a non-nil terminal error. After yielding a non-nil error, the sequence stops. Consumers should stop processing and return or break on the first non-nil error. On terminal errors, WithResult contains only lifecycle data observed before the error, and WithOnStats is not called.

Consumers using iter.Pull2 should prefer PullRowIteratorSeq, which normalizes terminal errors so pull returns ok=false when err!=nil.

func Rows

func Rows(rows ...*spanner.Row) iter.Seq2[*spanner.Row, error]

Rows adapts already-built rows to the fallible sequence shape used by RowIteratorSeq. Non-nil rows are yielded with a nil error. A nil row aborts the sequence by yielding ErrNilRow.

Row sources that can fail per row should produce their own iter.Seq2 instead of pre-building a slice for Rows.

func SliceToRowSeq

func SliceToRowSeq(rows []*spanner.Row) iter.Seq2[*spanner.Row, error]

SliceToRowSeq adapts an existing row slice to the fallible sequence shape used by RowIteratorSeq.

It exists for downstream tests, fakes, and virtual result sets that naturally store fixtures as []*spanner.Row. It is equivalent to Rows(rows...), including nil-row handling: nil rows yield ErrNilRow and abort the sequence.

Types

type Option

type Option func(*config)

Option configures RowIteratorSeq and DrainRowIterator.

func WithDrainOnEarlyStop

func WithDrainOnEarlyStop() Option

WithDrainOnEarlyStop configures RowIteratorSeq to consume the remaining rows after the consumer stops early.

Draining is disabled by default to preserve normal iterator early-exit behavior. Use this option when callers need WithOnStats to run after any early stop, including a range-loop break or an adapter that stops pulling because a downstream operation failed. Errors encountered only during this post-stop drain cannot be yielded to the caller and therefore suppress the stats hook. It has no effect on DrainRowIterator, which always drains.

func WithOnMetadata

func WithOnMetadata(f func(*sppb.ResultSetMetadata)) Option

WithOnMetadata registers a hook that runs once when result metadata becomes available.

For a query with rows, the hook runs after the first successful Next call and before that first row is yielded, so metadata captured by the hook is visible inside the first loop body. For an empty result set, the hook runs after Next returns iterator.Done. A nil hook is ignored.

func WithOnStats

func WithOnStats(f func(Stats)) Option

WithOnStats registers a hook that runs after the adapted iterator has reached iterator.Done and has been stopped.

If the consumer stops early, stats are available only when WithDrainOnEarlyStop is also configured. A nil hook is ignored.

func WithResult

func WithResult(result *RowIteratorResult) Option

WithResult stores iterator lifecycle data in result as it becomes available.

The pointed value is reset when iteration starts. Metadata is set before the first row is yielded, RowsRead is updated after each consumed row, and Stats is set only after the iterator reaches iterator.Done. On errors, result contains the partial lifecycle data observed before the error. A nil result is ignored.

func WithStatsEncoding added in v0.3.0

func WithStatsEncoding(enc StatsEncoding) Option

WithStatsEncoding configures how RowIteratorResult.StatsProto encodes row counts for a drained iterator. The default is StatsEncodingDefault.

type RowIteratorResult

type RowIteratorResult struct {
	Metadata *sppb.ResultSetMetadata
	Stats    Stats
	RowsRead int64
	// contains filtered or unexported fields
}

RowIteratorResult is the metadata and stats available from a cloud.google.com/go/spanner.RowIterator.

RowsRead counts rows consumed from the iterator. Metadata and Stats values are not deep-copied from the underlying RowIterator; treat returned maps and protos as read-only. Stats protobuf encoding is configured by WithStatsEncoding on the drain options.

Values returned or captured by RowIteratorSeq, DrainRowIterator, and PullRowIteratorSeq carry unexported lifecycle state required by RowIteratorResult.StatsProto and RowIteratorResult.ResultSet. Do not construct RowIteratorResult manually for protobuf conversion; use WithResult or the value returned from DrainRowIterator.

func DrainRowIterator

func DrainRowIterator(rowIter *spanner.RowIterator, opts ...Option) (*RowIteratorResult, error)

DrainRowIterator consumes rowIter to iterator.Done without yielding rows.

The helper owns rowIter and always calls *cloud.google.com/go/spanner.RowIterator.Stop before returning. It is useful when callers need result metadata, query stats, query plan, or DML row count but do not want to expose row values to application code. If iteration fails, the returned result can be non-nil and contain partial metadata and RowsRead observed before the error; stats are only populated after a successful drain to iterator.Done.

Cloud Spanner only populates metadata after the first Next call, and stats after Next returns iterator.Done. DrainRowIterator therefore still consumes the result stream internally; it does not ask Spanner for stats without reading the stream. To avoid reading data rows at the query level, callers must execute a statement that returns no data rows.

func (RowIteratorResult) ResultSet added in v0.3.0

func (r RowIteratorResult) ResultSet(rows []*structpb.ListValue) (*sppb.ResultSet, error)

ResultSet builds a protobuf ResultSet from materialized rows and iterator lifecycle data captured while draining a RowIterator.

rows may be nil when row values are intentionally omitted. Stats encoding comes from WithStatsEncoding on the drain options. Use only on package-produced RowIteratorResult values; see RowIteratorResult.StatsProto.

func (RowIteratorResult) StatsProto added in v0.3.0

func (r RowIteratorResult) StatsProto() (*sppb.ResultSetStats, error)

StatsProto returns captured stats as *sppb.ResultSetStats using the encoding configured by WithStatsEncoding when the iterator was drained.

Row counts are encoded only after stats were captured at iterator.Done. Query plan and query stats encode from the captured Stats value whenever present. Partial results from errors omit row counts even when StatsEncodingDMLExact is configured.

Call this only on RowIteratorResult values produced or filled by spaniter (WithResult, DrainRowIterator, PullRowIteratorSeq). Manual struct literals omit unexported lifecycle state and may omit row counts.

type Stats

type Stats struct {
	QueryPlan  *sppb.QueryPlan
	QueryStats map[string]any
	RowCount   int64
}

Stats holds execution information populated on a cloud.google.com/go/spanner.RowIterator after the iterator reaches iterator.Done.

QueryPlan and QueryStats are set when the query used QueryWithStats. RowCount holds the DML row count after iterator.Done. Values are not deep-copied from the underlying RowIterator; treat returned maps and protos as read-only.

Stats mirrors the public fields exposed by cloud.google.com/go/spanner.RowIterator. Use RowIteratorResult.StatsProto when downstream code needs the protobuf ResultSetStats shape. The Go client has already decoded query stats to a map and exposes row count as a single int64, so Stats cannot distinguish an absent row count from row_count_exact:0. The Go client's PartitionedUpdate APIs return counts directly rather than through RowIterator, so partitioned DML counts are outside this type's normal scope.

type StatsEncoding added in v0.3.0

type StatsEncoding int

StatsEncoding selects how captured stats are converted to protobuf ResultSetStats when using RowIteratorResult.StatsProto or RowIteratorResult.ResultSet.

Set encoding with WithStatsEncoding when draining a RowIterator.

const (
	// StatsEncodingDefault uses default query semantics: omit row_count_exact when
	// RowCount is zero because absent row count and exact zero are indistinguishable
	// on [cloud.google.com/go/spanner.RowIterator].
	StatsEncodingDefault StatsEncoding = iota
	// StatsEncodingDMLExact always encodes RowCount as row_count_exact, including
	// zero. Use only when the caller knows Stats came from executed standard DML
	// (not PLAN and not read-only queries).
	StatsEncodingDMLExact
)

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL