- Public Interfaces
- Proposed Changes
- Compatibility, Deprecation, and Migration Plan
- Test Plan
- Rejected Alternatives
Current State: Draft
Discussion Thread: link
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Kafka currently supports setting an upper bound on the number of requests allowed into the (incoming) request queue. This is an indirect way of controlling memory consumption and has a few drawbacks:
- An administrator needs to estimate the average request size in order to provide a meaningful size limit.
- This size limit may need to be periodically updated as the workload changes.
- The server is still susceptible to a simultaneous batch of large requests exhausting the JVM memory (causing an OOM exception).
The third scenario actually occurred a few times at LinkedIn - a sudden spike of very large request batches (1000s requests each) from a Hadoop job caused OOM exceptions on a production cluster.
This KIP proposes allowing an administrator to specify a memory limit in bytes, which resolves the above problems.
Also, the code/facilities introduced by this KIP might prove useful for similar problems on the client side - like producer/consumer memory bounds.
This KIP introduces a new server configuration parameter,
queued.max.request.bytes, that would specify a limit on the volume of requests that can be held in memory. This configuration parameter will co-exist with the existing
queued.max.requests (the code will respect both bounds and will not pick up new requests when either is hit).
Beyond the proposed new configuration key this KIP makes no changes to client or server public APIs.
- MemoryPoolAvgDepletedPercent - percent of the time request were not being read out of socket due to lack of memory
- MemoryPoolAvailable - number of bytes available in the pool
- MemoryPoolUsed - number of bytes currently allocated out of the pool and still not returned
a new MemoryPool interface would be introduced into the kafka codebase:
- the pool is non-blocking, so network threads would not be blocked waiting for memory and could make progress elsewhere.
- SocketServer would instantiate and hold a memory pool, which Processor threads would try to allocate memory from when reading requests out of sockets (by passing the pool to instances of NetworkReceive that they create).
- NetworkReceive.readFromReadableChannel() would be modified to try allocating memory (it is already written in a way that reading may involve multiple repeated calls to readFromReadableChannel(), so not a big change to behavior)
- memory would be released at the end of request processing (in KafkaRequestHandler.run()), and also in case of disconnection mid request-building in KafkaChannel.close()
- As the pool would allow any size request if it has any capacity available, the actual memory bound is
queued.max.request.bytes + socket.request.max.bytes. The up-side is no issues with large requests getting starved out
Request throttling by way of channel muting
when memory is unavailable, Selector would mute incoming channels so as not to get into a tight-loop where read-ready keys are returned by poll() but cannot actually be read from (because no memory). care must be taken when muting/unmuting because:
- SSL transports cannot be muted mid-handshake (I believe this is more for implementation simplicity in kafka than a fundamental limitation, but not one this KIP should address)
- muting a channel that has already been allocated memory to read into might result in a deadlock.
- muting channels is also kafka's way of ensuring only a single request is processed at a time from the same channel. this guarantee must be preserved.
- SSLTransportLayer has intermediate buffers where data may get "stuck" (due to no memory) and yet the underlying socket may be done and so will not show up on further select() calls.
so, the following changes are proposed in support of channel muting under memory pressure:
- TransportLayer.isInMutableState() would be introduced to account for transports that cannot currently be muted (currently only SSL during handshake)
- TransportLayer.hasBytesBuffered() would be introduced to account for transports that have unread data in any intermediate buffers (currently only possible for ssl transport)
- KafkaChannel.isInMutableState() would be introduced to account for channels with already allocated memory or channels who's underlying transport is not in a mutable state.
- Selector.poll() would mute all channels that are in a mutable state if the server has no memory to accept any further requests - this would prevent their keys from being returned by the underlying select() call and thus prevent a tight loop in SocketServer.Processor.run()
- when memory becomes available again (in some subsequent Selector.poll() call) any channels previously muted except those muted for single-request-procedding reasons (see #3 above) will be unmuted. this would require maintainign a set of channels explicitely muted (so they would not be unmuted when memory is available) in Selector
- channels who's underlying transports have yet-unread data buffered must be accounted for to prevent the case of a stale tail of data getting stuck in said buffers.
- concerns have been raised about the possibility of starvation in Selector.pollSelectionKeys() - in case the order of keys in Set<SelectionKey> selectionKeys is deterministic and memory is tight, sockets consistently at the beginning of the set get better treatment then those at the end of the iteration order. to overcome this code has been put in place to shuffle the selection keys and handle them in random order ONLY IF MEMORY IS TIGHT (so if the previous allocation call failed). this avoids the overhead of the shuffle when memory is not an issue.
Compatibility, Deprecation, and Migration Plan
There are a few approaches w.r.t migration. The current preference is to go with the third option.
queued.max.requestsis deprecated/removed in favor of
queued.max.request.bytes. In this case, the conversion of existing configurations could use
queued.max.request.bytes = queued.max.requests * socket.request.max.bytes(which is conservative, but "safe")
queued.max.requestsis supported as an alternative to
queued.max.request.bytes(either-or), in which case no migration is required. A default value of 0 could be used to disable the feature (by default) and runtime code would pick a queue implementation depending on which configuration parameter is provided.
queued.max.requestsis supported in addition
queued.max.request.bytes(both respected at the same time). In this case a default value of
queued.max.request.bytes = -1would maintain backwards compatible behavior.
The current naming scheme of
queued.max.requests (and the proposed
queued.max.request.bytes) may be a bit opaque. Perhaps using requestQueue.max.requests and requestQueue.max.bytes would more clearly convey the meaning to users (indicating that these settings deal with the request queue specifically, and not some other). The current
queued.max.requests configuration can be retained for a few more releases for backwards compatibility.
queued.max.request.bytes must be larger than socket.request.max.bytes (in other words, memory pool must be large enough to accommodate the largest single request possible), or <=0 (if disabled). the default would be -1.
- A unit test was written to validate the behavior of the memory pool
- A unit test that validates correct behavior of RequestChannel under capacity bounds would need to be written.
- A micro-benchmark for determining the performance of the pool would need to be written
- Stress testing a broker (heavy producer load of varying request sizes) to verify that the memory limit is honored.
- Benchmarking producer and consumer throughput before/after the change to prove that ingress/egress performance remains acceptable.
- Testing of SSL connections (both inter-broker and client-broker) since implementation bugs may affect ssl
- Reducing producer max batch size: this is harmful to throughput (and is also more complicated to maintain from an administrator's standpoint than simply sizing the broker itself). This is more of a workaround than a fix
- Reducing producer max request size: same issues as above.
- Limiting the number of connected clients: same issues as above
queued.max.requestsin the broker: Although this will conservatively size the queue it can be detrimental to throughput in the average case.
- controlling the volume of requests enqueued in
RequestChannel.requestQueue(would not suffice as no bound on memory read from actual sockets)
- the order of selection keys returned from a selector.poll call is undefined. in case the actual implementation uses a fixed order (say by increasing handle id?) and under prolonged memory pressure (so never enough memory to service all requests) this may lead to starvation of sockets that are always at the end of the iteration order. to overcome this the code shuffles the selection keys if memory is low.
- a strict pool (which adheres to its max size completely) will cause starvation of large requests under memory pressure (as they would never be able to allocate if there is a stream of small requests). to avoid this the pool implementation will allocate the requested amount of memory if it has any memory available (so if pool has 1 free byte and 1 MB is requested, 1MB will be returned and the number of available bytes in the pool will be negative). this means the actual bound on number of bytes outstanding is queued.max.request.bytes + socket.request.max.bytes - 1 (socket.request.max.bytes representing the largest single request possible)