[ 
https://issues.apache.org/jira/browse/FLINK-40944?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40944:
-----------------------------------
    Labels: pull-request-available  (was: )

>  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
>            Priority: Major
>              Labels: pull-request-available
>
> 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