RxJava Operator that dynamically buffers backpressured elements and emits them in batches












1















I have a Flowable that emits events that need to be handled by an expensive operation which expects element arrays:



Flowable<T> src
void expensiveOp(List<T> batch)


Other than using a constant window i'd like to specify a window of max elements that is filled while downstream is busy and when full just backpressures:



int maxSize = 1024
src.dynamicWindow(maxSize).subscribe(expensiveOp)


The size of the window should therefore be neither constant-time nor element but backpressure dependent. The buffer should be flushed when the subscriber is ready to process the next element.



What overloaded method am I missing?



Possible extensions would be a minSize parameter and a retry mechanism that retries with an increased window.










share|improve this question




















  • 2





    Sounds like you need coalesce

    – akarnokd
    Nov 15 '18 at 7:32
















1















I have a Flowable that emits events that need to be handled by an expensive operation which expects element arrays:



Flowable<T> src
void expensiveOp(List<T> batch)


Other than using a constant window i'd like to specify a window of max elements that is filled while downstream is busy and when full just backpressures:



int maxSize = 1024
src.dynamicWindow(maxSize).subscribe(expensiveOp)


The size of the window should therefore be neither constant-time nor element but backpressure dependent. The buffer should be flushed when the subscriber is ready to process the next element.



What overloaded method am I missing?



Possible extensions would be a minSize parameter and a retry mechanism that retries with an increased window.










share|improve this question




















  • 2





    Sounds like you need coalesce

    – akarnokd
    Nov 15 '18 at 7:32














1












1








1








I have a Flowable that emits events that need to be handled by an expensive operation which expects element arrays:



Flowable<T> src
void expensiveOp(List<T> batch)


Other than using a constant window i'd like to specify a window of max elements that is filled while downstream is busy and when full just backpressures:



int maxSize = 1024
src.dynamicWindow(maxSize).subscribe(expensiveOp)


The size of the window should therefore be neither constant-time nor element but backpressure dependent. The buffer should be flushed when the subscriber is ready to process the next element.



What overloaded method am I missing?



Possible extensions would be a minSize parameter and a retry mechanism that retries with an increased window.










share|improve this question
















I have a Flowable that emits events that need to be handled by an expensive operation which expects element arrays:



Flowable<T> src
void expensiveOp(List<T> batch)


Other than using a constant window i'd like to specify a window of max elements that is filled while downstream is busy and when full just backpressures:



int maxSize = 1024
src.dynamicWindow(maxSize).subscribe(expensiveOp)


The size of the window should therefore be neither constant-time nor element but backpressure dependent. The buffer should be flushed when the subscriber is ready to process the next element.



What overloaded method am I missing?



Possible extensions would be a minSize parameter and a retry mechanism that retries with an increased window.







rx-java rx-java2 reactivex






share|improve this question















share|improve this question













share|improve this question




share|improve this question








edited Dec 8 '18 at 18:30









Filip

3,23931226




3,23931226










asked Nov 14 '18 at 15:09









André RüdigerAndré Rüdiger

56111




56111








  • 2





    Sounds like you need coalesce

    – akarnokd
    Nov 15 '18 at 7:32














  • 2





    Sounds like you need coalesce

    – akarnokd
    Nov 15 '18 at 7:32








2




2





Sounds like you need coalesce

– akarnokd
Nov 15 '18 at 7:32





Sounds like you need coalesce

– akarnokd
Nov 15 '18 at 7:32












0






active

oldest

votes











Your Answer






StackExchange.ifUsing("editor", function () {
StackExchange.using("externalEditor", function () {
StackExchange.using("snippets", function () {
StackExchange.snippets.init();
});
});
}, "code-snippets");

StackExchange.ready(function() {
var channelOptions = {
tags: "".split(" "),
id: "1"
};
initTagRenderer("".split(" "), "".split(" "), channelOptions);

StackExchange.using("externalEditor", function() {
// Have to fire editor after snippets, if snippets enabled
if (StackExchange.settings.snippets.snippetsEnabled) {
StackExchange.using("snippets", function() {
createEditor();
});
}
else {
createEditor();
}
});

function createEditor() {
StackExchange.prepareEditor({
heartbeatType: 'answer',
autoActivateHeartbeat: false,
convertImagesToLinks: true,
noModals: true,
showLowRepImageUploadWarning: true,
reputationToPostImages: 10,
bindNavPrevention: true,
postfix: "",
imageUploader: {
brandingHtml: "Powered by u003ca class="icon-imgur-white" href="https://imgur.com/"u003eu003c/au003e",
contentPolicyHtml: "User contributions licensed under u003ca href="https://creativecommons.org/licenses/by-sa/3.0/"u003ecc by-sa 3.0 with attribution requiredu003c/au003e u003ca href="https://stackoverflow.com/legal/content-policy"u003e(content policy)u003c/au003e",
allowUrls: true
},
onDemand: true,
discardSelector: ".discard-answer"
,immediatelyShowMarkdownHelp:true
});


}
});














draft saved

draft discarded


















StackExchange.ready(
function () {
StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f53303267%2frxjava-operator-that-dynamically-buffers-backpressured-elements-and-emits-them-i%23new-answer', 'question_page');
}
);

Post as a guest















Required, but never shown

























0






active

oldest

votes








0






active

oldest

votes









active

oldest

votes






active

oldest

votes
















draft saved

draft discarded




















































Thanks for contributing an answer to Stack Overflow!


  • Please be sure to answer the question. Provide details and share your research!

But avoid



  • Asking for help, clarification, or responding to other answers.

  • Making statements based on opinion; back them up with references or personal experience.


To learn more, see our tips on writing great answers.




draft saved


draft discarded














StackExchange.ready(
function () {
StackExchange.openid.initPostLogin('.new-post-login', 'https%3a%2f%2fstackoverflow.com%2fquestions%2f53303267%2frxjava-operator-that-dynamically-buffers-backpressured-elements-and-emits-them-i%23new-answer', 'question_page');
}
);

Post as a guest















Required, but never shown





















































Required, but never shown














Required, but never shown












Required, but never shown







Required, but never shown

































Required, but never shown














Required, but never shown












Required, but never shown







Required, but never shown







Popular posts from this blog

Bressuire

Vorschmack

Quarantine