Skip to content

[SpannerIO] Add withWriteResults() to SpannerIO.Write to emit commit metadata - #40480

Open
stankiewicz wants to merge 1 commit into
apache:masterfrom
stankiewicz:spanner_io_write_results
Open

stankiewicz wants to merge 1 commit into
apache:masterfrom
stankiewicz:spanner_io_write_results

Conversation

@stankiewicz

Copy link
Copy Markdown
Contributor

Summary

Adds withWriteResults() to SpannerIO.Write and SpannerIO.WriteGrouped (modeled after BigtableIO.Write#withWriteResults()) to optionally enable Options.commitStats() on Cloud Spanner mutation writes and return a
PCollection<WriteResults> containing commit metadata from com.google.cloud.spanner.CommitResponse.

Changes

  • WriteResults:
    • Added WriteResults annotated with @DefaultCoder(AvroCoder.class) (using TimestampEncoding for Timestamp fields) exposing:
      • getMutationCount() (long): populated from CommitResponse.getCommitStats().getMutationCount().
      • getCommitTimestamp() (@Nullable Timestamp): populated from CommitResponse.getCommitTimestamp().
      • getSnapshotTimestamp() (@Nullable Timestamp): populated from CommitResponse.getSnapshotTimestamp().
  • SpannerIO.Write#withWriteResults() & SpannerIO.WriteGrouped#withWriteResults():
    • Added WriteWithResults (PTransform<PCollection<Mutation>, PCollection<WriteResults>>) and WriteGroupedWithResults (PTransform<PCollection<MutationGroup>, PCollection<WriteResults>>).
    • Enabled Options.commitStats() in WriteToSpannerFn#getTransactionOptions() only when withWriteResults() is used, preserving full backwards compatibility and zero overhead for existing SpannerIO.write() users.
    • Exposed @Nullable PCollection<WriteResults> getWriteResults() on SpannerWriteResult.
  • Bundle Minimum Timestamp in GatherSortCreateBatchesFn:
    • Updated GatherSortCreateBatchesFn and OutputReceiverForFinishBundle to track the minimum element timestamp in the bundle (minBundleTimestamp) and assign it to batches flushed at @FinishBundle instead of Instant. now().
  • Tests:
    • Added unit tests in SpannerIOWriteTest covering WriteResults default AvroCoder, minimum bundle timestamp in GatherSortCreateBatchesFn, Write#withWriteResults(), WriteGrouped#withWriteResults(), and timestamp
      preservation when batching is disabled.
    • Added integration test testWriteWithResults in SpannerWriteIT.

@github-actions

github-actions Bot commented Oct 9, 2026

Copy link
Copy Markdown
Contributor

Assigning reviewers:

R: @ahmedabu98 for label java.
R: @nielm for label spanner.

Note: If you would like to opt out of this review, comment assign to next reviewer.

Available commands:

  • stop reviewer notifications - opt out of the automated review tooling
  • remind me after tests pass - tag the comment author after tests pass
  • waiting on author - shift the attention set back to the author (any comment or push by the author will return the attention set to the reviewers)

The PR bot will only process comments in the main thread (not review comments).

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant