This is an automated email from the ASF dual-hosted git repository.
tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 2448f5ac7d6 Expand Beam Python heap dump with process/native memory
stats (#39244) (#39466)
2448f5ac7d6 is described below
commit 2448f5ac7d6a3bcbc61a7acd6edeeb400829ced1
Author: Udit Jain <[email protected]>
AuthorDate: Sat Sep 19 04:02:17 2026 +0530
Expand Beam Python heap dump with process/native memory stats (#39244)
(#39466)
The heap dump (--experiments=enable_heap_dump) previously only contained the
guppy Python-object heap. This adds a memory_stats() section so operators
can
tell whether memory is growing on the native (C) heap vs in Python objects,
and can spot native-heap fragmentation:
* Process memory: peak RSS (ru_maxrss, unit-corrected per platform) and, on
Linux, current RSS from /proc/self/statm.
* Python allocations: sys.getallocatedblocks() and per-generation gc stats.
* glibc malloc: mallinfo2() arena/hblkhd/uordblks/fordblks/keepcost plus a
fordblks/(arena+hblkhd) fragmentation ratio. mallinfo2 is used rather than
the legacy mallinfo, whose int fields overflow past 2GB and would
misreport
exactly when memory is the concern.
Every collector degrades gracefully and never raises when a source is
unavailable on the current platform. memory_stats() is emitted even when
guppy is not importable.
Tests cover section presence, that stats are emitted without guppy, and that
the glibc section degrades gracefully on non-glibc platforms.
Generated-by: Claude (Anthropic AI assistant)
Co-authored-by: uditjainstjis <[email protected]>
Co-authored-by: tvalentyn <[email protected]>
---
CHANGES.md | 2 +-
.../apache_beam/runners/worker/worker_status.py | 132 ++++++++++++++++++++-
.../runners/worker/worker_status_test.py | 25 ++++
3 files changed, 157 insertions(+), 2 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index cd4ea53eb76..5480ddc1c21 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -68,7 +68,7 @@
## New Features / Improvements
-* X feature added (Java/Python)
([#X](https://github.com/apache/beam/issues/X)).
+* (Python) Expanded the SDK worker heap dump
(`--experiments=enable_heap_dump`) with process RSS, CPython allocator/GC
stats, and glibc `mallinfo2` native-heap/fragmentation stats to help
distinguish native-heap from Python-object memory growth
([#39244](https://github.com/apache/beam/issues/39244)).
## Breaking Changes
diff --git a/sdks/python/apache_beam/runners/worker/worker_status.py
b/sdks/python/apache_beam/runners/worker/worker_status.py
index 1d54a3ee176..28917c6d949 100644
--- a/sdks/python/apache_beam/runners/worker/worker_status.py
+++ b/sdks/python/apache_beam/runners/worker/worker_status.py
@@ -19,6 +19,7 @@
import gc
import logging
+import os
import queue
import sys
import threading
@@ -108,6 +109,135 @@ def thread_dump(thread_prefix=None):
return '\n'.join(all_traces)
+def _process_memory_stats():
+ """Process-level memory usage (RSS), to gauge total footprint over time."""
+ lines = ['--- Process memory ---']
+ try:
+ import resource
+ usage = resource.getrusage(resource.RUSAGE_SELF)
+ # ru_maxrss is reported in kilobytes on Linux but in bytes on macOS/BSD.
+ if sys.platform == 'darwin':
+ peak_bytes = usage.ru_maxrss
+ else:
+ peak_bytes = usage.ru_maxrss * 1024
+ lines.append('peak RSS (ru_maxrss): %d bytes' % peak_bytes)
+ except Exception as e: # pylint: disable=broad-except
+ lines.append('resource stats unavailable: %s' % e)
+ try:
+ # /proc/self/statm reports counts in pages; field 1 is resident set size.
+ with open('/proc/self/statm') as f:
+ resident_pages = int(f.read().split()[1])
+ lines.append(
+ 'current RSS (/proc/self/statm): %d bytes' %
+ (resident_pages * os.sysconf('SC_PAGE_SIZE')))
+ except Exception: # pylint: disable=broad-except
+ # /proc is Linux-only; skip silently on other platforms.
+ pass
+ return '\n'.join(lines)
+
+
+def _python_memory_stats():
+ """CPython allocator and garbage-collector stats.
+
+ Together with the native-heap stats these help distinguish memory growth in
+ Python objects from growth on the native (C) heap.
+ """
+ lines = ['--- Python allocations ---']
+ try:
+ lines.append('sys.getallocatedblocks: %d' % sys.getallocatedblocks())
+ except Exception as e: # pylint: disable=broad-except
+ lines.append('sys.getallocatedblocks unavailable: %s' % e)
+ try:
+ lines.append('gc.get_count (gen0, gen1, gen2): %s' % (gc.get_count(), ))
+ for i, stat in enumerate(gc.get_stats()):
+ lines.append(
+ 'gc gen%d: collections=%s collected=%s uncollectable=%s' % (
+ i,
+ stat.get('collections'),
+ stat.get('collected'),
+ stat.get('uncollectable')))
+ except Exception as e: # pylint: disable=broad-except
+ lines.append('gc stats unavailable: %s' % e)
+ return '\n'.join(lines)
+
+
+def _glibc_malloc_stats():
+ """glibc allocator stats, useful for reasoning about native-heap growth and
+ fragmentation.
+
+ ``fordblks`` (free space the allocator retains rather than returning to the
+ OS) growing relative to ``uordblks`` (space in use) is a signal of native
+ heap fragmentation. Only available with glibc (Linux); degrades gracefully
+ elsewhere.
+ """
+ lines = ['--- glibc malloc (native heap) ---']
+ if not sys.platform.startswith('linux'):
+ lines.append('unavailable: glibc malloc stats are only collected on
Linux.')
+ return '\n'.join(lines)
+ try:
+ import ctypes
+
+ class _MallInfo2(ctypes.Structure):
+ # Mirrors glibc's ``struct mallinfo2`` (all fields are size_t).
+ _fields_ = [
+ ('arena', ctypes.c_size_t),
+ ('ordblks', ctypes.c_size_t),
+ ('smblks', ctypes.c_size_t),
+ ('hblks', ctypes.c_size_t),
+ ('hblkhd', ctypes.c_size_t),
+ ('usmblks', ctypes.c_size_t),
+ ('fsmblks', ctypes.c_size_t),
+ ('uordblks', ctypes.c_size_t),
+ ('fordblks', ctypes.c_size_t),
+ ('keepcost', ctypes.c_size_t),
+ ]
+
+ libc = ctypes.CDLL('libc.so.6')
+ if not hasattr(libc, 'mallinfo2'):
+ # The older mallinfo() uses int fields that overflow past 2GB and would
+ # misreport exactly when memory is the concern, so we do not fall back.
+ lines.append('unavailable: mallinfo2 not found (needs glibc >= 2.33).')
+ return '\n'.join(lines)
+ libc.mallinfo2.restype = _MallInfo2
+ libc.mallinfo2.argtypes = []
+ info = libc.mallinfo2()
+ lines.append('arena (non-mmapped bytes from sbrk): %d' % info.arena)
+ lines.append('hblkhd (mmapped bytes): %d' % info.hblkhd)
+ lines.append('uordblks (in-use bytes): %d' % info.uordblks)
+ lines.append(
+ 'fordblks (free bytes retained by allocator): %d' % info.fordblks)
+ lines.append(
+ 'keepcost (releasable top-most free bytes): %d' % info.keepcost)
+ footprint = info.arena + info.hblkhd
+ if footprint:
+ lines.append(
+ 'fragmentation (fordblks / (arena + hblkhd)): %.2f%%' %
+ (100.0 * info.fordblks / footprint))
+ except Exception as e: # pylint: disable=broad-except
+ lines.append('unavailable: %s' % e)
+ return '\n'.join(lines)
+
+
+def memory_stats():
+ """Collect process, Python and native-heap memory statistics.
+
+ This complements the guppy heap dump with information that helps distinguish
+ memory growth on the native (C) heap from Python-object allocations, and with
+ glibc allocator stats that hint at native-heap fragmentation. Every collector
+ degrades gracefully and never raises when a source is unavailable on the
+ current platform.
+ """
+ banner = '=' * 10 + ' MEMORY STATS ' + '=' * 10
+ sections = [
+ banner,
+ _process_memory_stats(),
+ _python_memory_stats(),
+ _glibc_malloc_stats(),
+ '=' * 30,
+ ]
+ return '\n'.join(sections)
+
+
def heap_dump():
"""Get a heap dump for the current SDK worker harness. """
banner = '=' * 10 + ' HEAP DUMP ' + '=' * 10 + '\n'
@@ -116,7 +246,7 @@ def heap_dump():
else:
heap = '%s\n' % hpy().heap()
ending = '=' * 30
- return banner + heap + ending
+ return banner + heap + ending + '\n' + memory_stats()
def _state_cache_stats(state_cache: StateCache) -> str:
diff --git a/sdks/python/apache_beam/runners/worker/worker_status_test.py
b/sdks/python/apache_beam/runners/worker/worker_status_test.py
index 88543258250..8ddd88813ac 100644
--- a/sdks/python/apache_beam/runners/worker/worker_status_test.py
+++ b/sdks/python/apache_beam/runners/worker/worker_status_test.py
@@ -28,6 +28,7 @@ from apache_beam.portability.api import beam_fn_api_pb2_grpc
from apache_beam.runners.worker import statesampler
from apache_beam.runners.worker.worker_status import FnApiWorkerStatusHandler
from apache_beam.runners.worker.worker_status import heap_dump
+from apache_beam.runners.worker.worker_status import memory_stats
from apache_beam.utils import thread_pool_executor
from apache_beam.utils.counters import CounterName
@@ -231,6 +232,30 @@ class HeapDumpTest(unittest.TestCase):
self.assertTrue(
'Unable to import guppy, the heap dump will be skipped' in result)
+ @mock.patch('apache_beam.runners.worker.worker_status.hpy', None)
+ def test_heap_dump_includes_memory_stats_without_guppy(self):
+ # Memory stats must be emitted even when guppy is unavailable.
+ result = '%s' % heap_dump()
+ self.assertIn('MEMORY STATS', result)
+
+
+class MemoryStatsTest(unittest.TestCase):
+ def test_memory_stats_sections_present(self):
+ result = memory_stats()
+ self.assertIn('MEMORY STATS', result)
+ self.assertIn('Process memory', result)
+ self.assertIn('Python allocations', result)
+ self.assertIn('glibc malloc', result)
+ # Python allocator introspection is available on every platform.
+ self.assertIn('sys.getallocatedblocks', result)
+
+ @mock.patch('sys.platform', 'darwin')
+ def test_glibc_stats_graceful_on_non_glibc(self):
+ # On non-glibc platforms the native-heap section must degrade gracefully
+ # rather than raise.
+ result = memory_stats()
+ self.assertIn('unavailable', result)
+
if __name__ == '__main__':
logging.getLogger().setLevel(logging.INFO)