João Boto created FLINK-40944:
---------------------------------
Summary: CPU limit-factor inflates the JVM's
availableProcessors() to the CPU limit
Key: FLINK-40944
URL: https://issues.apache.org/jira/browse/FLINK-40944
Project: Flink
Issue Type: Bug
Components: Kubernetes Operator
Reporter: João Boto
h3. Description
CPU limit-factor inflates the JVM's availableProcessors() to the CPU limit,
risking OutOfMemoryError: Direct buffer memory — operator should pin
-XX:ActiveProcessorCount to the CPU request
h3.
Problem
* kubernetes.\{jobmanager,taskmanager}.cpu.limit-factor (FLINK-15648) raises
the container CPU limit above the request (e.g. request 10, factor 4 → limit
40) to provide burst headroom. Nothing else in the memory model accounts for
this:
The container-aware JVM derives Runtime.availableProcessors() from the cgroup
CPU quota, which reflects the limit, not the request — so the JVM sees 40 cores
instead of 10.
* Thread pools sized off that count scale up accordingly — e.g. Netty's pooled
direct-memory allocator (PooledByteBufAllocator arenas/threads in Flink's
network stack), Beam portability/bundle-factory pools in PyFlink jobs.
* But MaxDirectMemorySize is computed from the fixed, request-sized
taskmanager.memory.* model and never changes with the CPU factor.
Net effect: ~factor× more thread-driven direct-buffer demand against an
unchanged direct-memory ceiling → java.lang.OutOfMemoryError: Direct buffer
memory.
h3.
Real-world impact
Production TaskManager with cpu: 10, cpu.limit-factor: 4 (limit 40),
taskmanager.memory.process.size: 20G repeatedly failed with OutOfMemoryError:
Direct buffer memory inside a Kafka source fetcher thread, right at the JVM's
~1.6 GiB direct-memory ceiling. Reproducible and directly tied to the inflated
processor count. Neither FLINK-15648, the operator's ResourceRequirements
rework, nor the fractional-CPU discussions on the user list mention this
interaction.
h3. Proposed solution
When the effective CPU limit-factor > 1, the operator should append
-XX:ActiveProcessorCount=<max(1, ceil(cpu request))> to
env.java.opts.jobmanager/env.java.opts.taskmanager in the effective
configuration (config map), so the JVM stays sized to the request regardless of
the inflated cgroup quota. This is a no-op when limit == request, and covers
both JM and TM in native and standalone modes since JVM opts flow from the
shared config. An explicit -XX:ActiveProcessorCount already present in the
user's env.java.opts.* must win (skip injection).
FlinkConfigBuilder is the natural place: it already computes the effective
limit factor (setResourceRequirements) and reads the deprecated-path values
from kubernetes.*.cpu.limit-factor (same options
FlinkUtils.calculateClusterCpuUsage consumes), so a single post-build step
covers both the new resources field and the raw-config path.
Precedent: Flink core already injects JVM options into env.java.opts for k8s
(KubernetesUtils.addGCLogsToJavaOpts for GC logging); Spark resolved the
identical issue by setting ActiveProcessorCount on driver/executor JVMs
(SPARK-31028).
Alternatives considered:
* Fix in flink-kubernetes core where the limit-factor is applied to the
container spec — broader (covers non-operator deployments too) but the JVM
command is assembled at entrypoint/process-launch time, making config-side
injection on the operator the cleaner first step.
* Documentation-only recommendation. Weaker: the failure mode is a silent prod
OOM that appears only when the feature is used as designed.
h3.
Workaround today (what we run)
Users can render -XX:ActiveProcessorCount=<ceil(cpu request)> via
env.java.opts.\{jobmanager,taskmanager} themselves whenever a limit-factor is
set — which is what our platform wrapper does; but this shouldn't be required
to use a first-class operator/feature safely.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)