Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.md
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@
* (Python) Removed the `envoy-data-plane` (and transitive `betterproto`) dependency; `EnvoyRateLimiter` now uses a small vendored protobuf definition instead, resolving dependency conflicts for downstream projects ([#37854](https://github.com/apache/beam/issues/37854)).
* (Java) Supported acknowledge mode for JmsIO ([#39253](https://github.com/apache/beam/issues/39253)).
* (Python) Added `equal_to_approx`, an `assert_that` matcher that compares numeric pipeline outputs with a configurable tolerance ([#18028](https://github.com/apache/beam/issues/18028)).
* (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)).
* X feature added (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).

## Breaking Changes
Expand Down
132 changes: 131 additions & 1 deletion sdks/python/apache_beam/runners/worker/worker_status.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@

import gc
import logging
import os
import queue
import sys
import threading
Expand Down Expand Up @@ -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'
Expand All @@ -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:
Expand Down
25 changes: 25 additions & 0 deletions sdks/python/apache_beam/runners/worker/worker_status_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
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

Expand Down Expand Up @@ -231,6 +232,30 @@ def test_skip_heap_dump(self):
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)
Expand Down
Loading