GitHub user uce opened a pull request:
https://github.com/apache/flink/pull/2883
[FLINK-5088] [network] Add capacity limit option to partitions
We worked with @StephanEwen on this. The changes build on top of FLINK-5159
(see #2882). Only the most recent commit is relevant.
# Changes
The goal of the changes was to allow limiting the capacity of produced
partitions in addition to the back pressure that is introduced via the buffer
pools.
## PIPELINED_BOUNDED partition type
This adds a further result partition type (`PIPELINED_BOUNDED`) that is
translated to pipelined runtime partitions that can configure a maximum
capacity for their output queues via:
taskmanager.net.max-out-queue-length
By default, they have unbounded length (`0`) and behave like regular
PIPELINED results.
## RemoteInputChannel capacity limit
If the outbound queues are limited, we need to capacity limit the inbound
remote channels, too. Otherwise, we only shift the problem to the receiver,
which potentially consume the bounded queues fast and buffers the data.
## Known Issues
For the time being the back pressure stats reported in the web interface
are not correct as they only take the buffer pools into account. If we keep the
default at unlimited, it should be OK to address this in a follow up.
You can merge this pull request into a Git repository by running:
$ git pull https://github.com/uce/flink 5088-bounded_queue
Alternatively you can review and apply these changes as the patch at:
https://github.com/apache/flink/pull/2883.patch
To close this pull request, make a commit to your master/trunk branch
with (at least) the following in the commit message:
This closes #2883
----
commit 2aad264e6d9273d15b0e04527f0f2593d2951e8c
Author: Stephan Ewen <[email protected]>
Date: 2016-11-27T17:15:40Z
[FLINK-5169] [network] Add tests for channel consumption
commit 212d6ac7c86a4426b6e4418ae5f9d1d4680ec05b
Author: Stephan Ewen <[email protected]>
Date: 2016-11-28T08:59:29Z
[FLINK-5169] [network] Make consumption of InputChannels fair
commit 564b96d72b1536f4ab171c723e3ae0ab2e3c935c
Author: Ufuk Celebi <[email protected]>
Date: 2016-11-28T08:59:58Z
[FLINK-5169] [network] Adjust tests to new consumer logic
commit 80246fac7fb1de6d24f17416e85435bcf0a0f1ad
Author: Ufuk Celebi <[email protected]>
Date: 2016-11-17T17:01:53Z
[FLINK-5088] [network] Add capacity limit option to partitions
This adds a further result partition type (PIPELINED_BOUNDED)
that is translated to pipelined runtime partitions that can
configure a maximum capacity for their output queues via:
taskmanager.net.max-out-queue-length
By default, they have unbounded length and behave like
regular PIPELINED results.
If the outbound queues are limited, we need to capacity limit
the inbound remote channels, too. Otherwise, we only shift the
problem to the receiver.
commit 644266ef1a6da767acdfd9823f7103248f970bfc
Author: Ufuk Celebi <[email protected]>
Date: 2016-11-28T10:57:48Z
Fix StreamRecordWriterTest
----
---
If your project is set up for it, you can reply to this email and have your
reply appear on GitHub as well. If your project does not have this feature
enabled and wishes so, or if the feature is enabled but not working, please
contact infrastructure at [email protected] or file a JIRA ticket
with INFRA.
---