-
Notifications
You must be signed in to change notification settings - Fork 4.3k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
[BEAM-10212] Integrate caching client #15214
[BEAM-10212] Integrate caching client #15214
Conversation
R: @amaliujia |
0669dea
to
5184b06
Compare
Codecov Report
@@ Coverage Diff @@
## master #15214 +/- ##
==========================================
- Coverage 83.82% 83.82% -0.01%
==========================================
Files 441 441
Lines 59706 59706
==========================================
- Hits 50051 50048 -3
- Misses 9655 9658 +3
Continue to review full report at Codecov.
|
Overall LGTM. Please verify if the ordering timer test works. |
460fac3
to
6c7f798
Compare
Run Java_Examples_Dataflow PreCommit |
Run Java_Examples_Dataflow_Java11 PreCommit |
// specified. If pipeline is batch, use a CachingBeamFnStateClient to store state responses. | ||
// User state caching is currently not supported in streaming mode. | ||
HandleStateCallsForBundle beamFnStateClient; | ||
if (bundleDescriptor.hasStateApiServiceDescriptor()) { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Can you explain a bit more why hasStateApiServiceDescriptor
is true = batch?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
hasStateApiServiceDescriptor
determines if a state handler is used at all, and options.as(StreamingOptions.class).isStreaming()
determines if it is batch or streaming. Updated comments to show this more clearly
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
O I see. options.as(StreamingOptions.class).isStreaming()
is used at line 530.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Enabled cache can also pass internal test suite.
Run Java_Examples_Dataflow_Java11 PreCommit |
Run Java_Examples_Dataflow PreCommit |
Run Java PreCommit |
2 similar comments
Run Java PreCommit |
Run Java PreCommit |
The failing |
// Instantiate a State API call handler depending on whether a State ApiServiceDescriptor was | ||
// specified. | ||
HandleStateCallsForBundle beamFnStateClient; | ||
if (bundleDescriptor.hasStateApiServiceDescriptor()) { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
@amaliujia @anthonyqzhu
I was under the impression that we would place the cache behind an experiment instead of opting everyone into its usage.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
We merged this assuming it was okay after seeing the passing tests / TAP presubmit passing as well.
I can introduce an experiment in front of this block today if necessary
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Please do.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Agreed. Let's file a patch to hide the cache by an experiment.
Adds state cache and caching client to process bundle execution
Thank you for your contribution! Follow this checklist to help us incorporate your contribution quickly and easily:
R: @username
).[BEAM-XXX] Fixes bug in ApproximateQuantiles
, where you replaceBEAM-XXX
with the appropriate JIRA issue, if applicable. This will automatically link the pull request to the issue.CHANGES.md
with noteworthy changes.See the Contributor Guide for more tips on how to make review process smoother.
ValidatesRunner
compliance status (on master branch)Examples testing status on various runners
Post-Commit SDK/Transform Integration Tests Status (on master branch)
Pre-Commit Tests Status (on master branch)
See .test-infra/jenkins/README for trigger phrase, status and link of all Jenkins jobs.
GitHub Actions Tests Status (on master branch)
See CI.md for more information about GitHub Actions CI.