How to manage cpu and memory with django channels

Viewed 540

Is there a way to better manage cpu and memory usage with django channels? I am trying to stream data very quickly over websockets, which works well until my ubuntu server cpu hits 100% and memory hits 90%, then daphne crashes. I am using NGINX and Daphne as server/proxy.

Is there a better way to setup redis? or potentially a method to 'clear cache' or something similar?

my setup is:

settings.py

CHANNEL_LAYERS = {
    "default": {
        "BACKEND": "channels_redis.core.RedisChannelLayer",
        "CONFIG": {
            "hosts": [("localhost", 6379)],
            "capacity": 5000,
            "expiry": 5,
        },
    },
}

asgi.py

application = ProtocolTypeRouter({
    # Django's ASGI application to handle traditional HTTP requests
    "http": django_asgi_app,

    # WebSocket chat handler
    "websocket": AuthMiddlewareStack(
        URLRouter([
            url(r"^stream/(?P<device_id>[\d\-]+)/$", StreamConsumer.as_asgi()),
            url(r"^status/(?P<device_id>[\d\-]+)/$", StatusConsumer.as_asgi())
        ])
    ),
})

consumer.py

class StreamConsumer(AsyncWebsocketConsumer):
    async def connect(self):
        self.room_name = self.scope['url_route']['kwargs']['device_id']
        self.room_group_name = 'stream_%s' % self.room_name
        self.token_passed = self.scope['query_string'].decode() #token passed in with ws://url/?token
        self.token_actual = await self.get_token(self.room_name) #go get actual token for this monitor

        # Join room group
        await self.channel_layer.group_add(
            self.room_group_name,
            self.channel_name
        )

        #check to ensure tokens match, if they do, allow connection
        if self.token_passed == self.token_actual:
            await self.accept()
            print('True')
        else:
            print(self.token_passed)
            print(self.token_actual)
            await self.close()


    async def disconnect(self, close_code):
        # Leave room group
        print(close_code)
        print("closing")
        await self.channel_layer.group_discard(
            self.room_group_name,
            self.channel_name
        )

    @database_sync_to_async
    def get_token(self, monitor_serial):
        monitor = CustomUser.objects.get(deviceSerial=monitor_serial)
        token = Token.objects.get(user=monitor)
        return token.key


    # Receive message from WebSocket
    async def receive(self, bytes_data=None, text_data=None):
        # Send message to room group
        await self.channel_layer.group_send(
            self.room_group_name,
            {
                'type': 'stream_data',
                'bytes_data': bytes_data,
                'text_data': text_data
            }
        )

    # Receive message from room group
    async def stream_data(self, event):
        bytes_data = event['bytes_data']
        text_data = event['text_data']
        # Send message to WebSocket
        await self.send(bytes_data=bytes_data, text_data=text_data)

python Code sending data to WebSocket:

#!/usr/bin/python

import websocket
import time

def on_message(ws, message):
    try:
        if message == 'pong':
           global pong
           global receiver_exists
           pong = True
           receiver_exists = True
    except:
        pass

def on_error(ws, error):
    print("### error ###", error)

def on_close(ws):
    print("### closed ###")

def on_open(ws):
    print('### connected ###')
    def stream(*args):
        while True:
            ws.send('large string')
            time.sleep(0.0005)

    stream()


while True:
    uri = f"mysocketurl"
    ws = websocket.WebSocketApp(uri,
                        on_open = on_open,
                        on_message = on_message,
                        on_error = on_error,
                        on_close = on_close)

    ws.run_forever()

The error message I am getting then daphne crashes is:

2021-05-08 13:45:14,398 WARNING  Application instance <Task pending coro=<ProtocolTypeRouter.__call__() running at /home/polysense/polysensesite/psvirtualenv/lib/python3.6/site-packages/channels/routing.py:71> wait_for=<Future pending cb=[<TaskWakeupMethWrapper object at 0x7efcf0ab42e8>()]>> for connection <WebSocketProtocol client=['127.0.0.1', 58056] path=b'/stream/1927-0000-0001/'> took too long to shut down and was killed.

You might also ask, why do I need to send data so quickly. I am trying to send jpg data using open-cv to create a video stream. I thought it was open-cv causing the crash, but it appears to just be trying to send the jpg data too quickly. When I slow it down by putting a delay of 100ms it will not crash, but the frame rate is not acceptable.

0 Answers
Related