nbali opened a new pull request, #40391:
URL: https://github.com/apache/beam/pull/40391
Addresses #27170 (alleviates it, does not fix it).
## What
FirestoreIO batch writes now grow their ramp-up budget from one
pipeline-wide start time, instead of from each DoFn instance's first write.
- `FirestoreV1.BatchWriteWithSummary` and `BatchWriteWithDeadLetterQueue`
create a singleton side input holding `clock.instant()`, computed once per
pipeline (`Create` → `MapElements` → `View.asSingleton()`), and pass it to the
write `ParDo`.
- `FirestoreV1WriteFn.BaseBatchWriteFn.processElement` reads the side input
and hands it to `RpcQos.setRampUpStart(Instant)`.
- `RpcQosImpl.WriteRampUp` uses that instant as its ramp-up start. If it's
never set, behavior is unchanged: the ramp-up starts lazily at the first write.
## Why
#27170 describes how slow the write ramp-up is for small jobs. On top of the
low starting budget (`500 / hintMaxNumWorkers`, 1 write/s with the default
hint), each `RpcQos` instance (one per DoFn instance, created in `@Setup`)
starts its own ramp-up clock at its first write. Every DoFn instance created
later starts again at the base budget. This happens on autoscaled workers, when
concurrency grows, and after a failed bundle discards its DoFn instances. Time
the pipeline spends before its first write doesn't count either.
With a shared start time:
- DoFn instances created later join at the budget the pipeline has already
grown to, instead of starting over at 1 write/s;
- time spent before the first write (for example a long read or transform
stage) counts toward the ramp-up.
Trade-off, the same one DatastoreIO accepts: if the first write happens long
after the pipeline starts, the budget has already grown, so those writes get
less ramp-up protection.
## Where it comes from
This copies the approach DatastoreIO already uses.
`DatastoreV1.Mutate.expand` creates a "Generate start timestamp" side input
(`Create.of("side input")` → `MapElements.via(s -> Instant.now())` →
`View.asSingleton()`) and passes it to `RampupThrottlingFn`. Its
`processElement` reads `c.sideInput(firstInstantSideInput)` and computes the
budget relative to that instant. This PR applies the same pattern to the
Firestore write path. The Firestore ramp-up lives inside `RpcQosImpl` rather
than in a separate DoFn, so the instant is passed into it.
## Pipeline update compatibility
The write transforms gain the side-input transforms, which changes the
pipeline graph. Following the existing pattern (for example `View`), the side
input is only added when `--updateCompatibilityVersion` is unset or at least
`2.78.0`. Pipelines that need to update a running job from an older SDK can set
it to an older version and keep the previous shape and behavior.
## Why this alleviates rather than fixes #27170
The rate formula is unchanged. A small job that writes from the start still
begins at `500 / hintMaxNumWorkers` writes per second and grows by 50% every 5
minutes, so with the default hint it is still slow for the first 10–15 minutes.
This PR removes the extra slowdown from per-instance restarts and makes
Firestore behave like Datastore. Changing the default growth is left for a
follow-up.
## Testing
- `RpcQosTest.rampUp_growsFromProvidedStart`: with a start 20 minutes before
the first write, the first batch is sized from the grown budget (3), not the
base budget (1).
- `FirestoreV1FnBatchWriteWithSummaryTest.rampUpStartIsReadFromSideInput`:
the DoFn reads the side input and passes it to `RpcQos.setRampUpStart`.
- `FirestoreV1BatchWriteRampUpStartTest`: both write transforms add the side
input to their `batchWrite` `ParDo`; with `updateCompatibilityVersion=2.77.0`
they don't.
------------------------
Thank you for your contribution! Follow this checklist to help us
incorporate your contribution quickly and easily:
- [x] Mention the appropriate issue in your description (for example:
`addresses #123`), if applicable. This will automatically add a link to the
pull request in the issue. If you would like the issue to automatically close
on merging the pull request, comment `fixes #<ISSUE NUMBER>` instead.
- [x] Update `CHANGES.md` with noteworthy changes.
- [ ] ~If this contribution is large, please file an Apache [Individual
Contributor License Agreement](https://www.apache.org/licenses/icla.pdf).~
See the [Contributor Guide](https://beam.apache.org/contribute) for more
tips on [how to make review process
smoother](https://github.com/apache/beam/blob/master/CONTRIBUTING.md#make-the-reviewers-job-easier).
To check the build health, please visit
[https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md](https://github.com/apache/beam/blob/master/.test-infra/BUILD_STATUS.md)
GitHub Actions Tests Status (on master branch)
------------------------------------------------------------------------------------------------
[](https://github.com/apache/beam/actions?query=workflow%3A%22Build+python+source+distribution+and+wheels%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Python+Tests%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Java+Tests%22+branch%3Amaster+event%3Aschedule)
[](https://github.com/apache/beam/actions?query=workflow%3A%22Go+tests%22+branch%3Amaster+event%3Aschedule)
See [CI.md](https://github.com/apache/beam/blob/master/CI.md) for more
information about GitHub Actions CI or the [workflows
README](https://github.com/apache/beam/blob/master/.github/workflows/README.md)
for a list of workflows and how to trigger them.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]