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)
   
------------------------------------------------------------------------------------------------
   [![Build python source distribution and 
wheels](https://github.com/apache/beam/actions/workflows/build_wheels.yml/badge.svg?event=schedule&&?branch=master)](https://github.com/apache/beam/actions?query=workflow%3A%22Build+python+source+distribution+and+wheels%22+branch%3Amaster+event%3Aschedule)
   [![Python 
tests](https://github.com/apache/beam/actions/workflows/python_tests.yml/badge.svg?event=schedule&&?branch=master)](https://github.com/apache/beam/actions?query=workflow%3A%22Python+Tests%22+branch%3Amaster+event%3Aschedule)
   [![Java 
tests](https://github.com/apache/beam/actions/workflows/java_tests.yml/badge.svg?event=schedule&&?branch=master)](https://github.com/apache/beam/actions?query=workflow%3A%22Java+Tests%22+branch%3Amaster+event%3Aschedule)
   [![Go 
tests](https://github.com/apache/beam/actions/workflows/go_tests.yml/badge.svg?event=schedule&&?branch=master)](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]

Reply via email to