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)

Reply via email to