Memory leak occurs when using Gstreamer and multiprocessing.Manager().Namespace() in the same process

Viewed 170

I'm using a Gstreamer to storyboard some video. I put each frame in multiprocessing.Manager().Namespace() for transfer between processes. In this case, a memory leak is observed. During 10 hours of application work, memory consumption increased by 400 MB. If you comment out the call of self.update_frame(buf.extract_dup(0, buf.get_size())), then the memory leak does not occur.

An example of a memory leak test:

import unittest
import gi
import traceback
import os, psutil
from multiprocessing import Process, Manager, Event

gi.require_version('Gst', '1.0')
from gi.repository import Gst


class RtpNamespaceTest(unittest.TestCase):
    pipeline_str = '''
        videotestsrc pattern=ball ! \
        appsink name=handle-app-sink \
        emit-signals=True \
        max-buffers=1 \
        drop=True \
        '''

    name_space = Manager().Namespace()
    pipeline = None
    event_interrupt: Event = Event()

    def start(self):
        # initializing gstreamer, subscribing to rtp-stream
        Gst.init(None)
        print(self.pipeline_str)
        self.pipeline = Gst.parse_launch(self.pipeline_str)
        self.pipeline.set_state(Gst.State.PLAYING)
        self.appsink = self.pipeline.get_by_name('handle-app-sink')
        self.appsink.connect("new-sample", self.on_new_buffer)

        bus = self.pipeline.get_bus()
        # listen to messages until storyboarding is stopped
        while not self.event_interrupt.is_set():
            bus.timed_pop_filtered(10000, Gst.MessageType.ANY)

        # stops the pipeline, free up resources
        self.pipeline.set_state(Gst.State.NULL)

    def on_new_buffer(self, src):
        sample = src.emit("pull-sample")
        buf = sample.get_buffer()
        # writing a frame in binary format in Namespace()
        self.update_frame(buf.extract_dup(0, buf.get_size()))
        print(f"RAM = {psutil.Process(os.getpid()).memory_info().rss / 1024 / 1024}")
        return Gst.FlowReturn.OK
        
    def update_frame(self, frame: bytes):
        self.name_space.frame = frame

    def test_repository(self):
        while True:
            try:
                self.start()
            except Exception as ex:
                traceback.print_exc()
                time.sleep(float(10))
0 Answers
Related