How to call a second fuction in the background while running the main program?

Viewed 44

I have one program that collects data from a websocket, processes the data and if some conditions apply I want to call another function that does something with the data.

This is easy enough, but I want the program that collects the data from the websocket to keep running.

I have 'fixed' this quite ugly by writing the data in a database and letting the second program check the database every few seconds. But I don't want to use this solution, since I occasionally get database is locked errors.

Is there a way to start program B from program A while program A keeps running? I have looked at multi threading and multi processing, and I feel this could be a way to solve it, but while I grasp the basic of that, it is still a bit too difficult for me to use.

Is there an easier way? and if not should I study multi threading or multi processing more?

(or if anyone knows a good guide/video, that would be great too!)

2 Answers

I suggest launching a worker thread, waiting for data to process. Main thread listen to websocket, and send data to worker through pipe.

The logic of worker is:

while True:
    data = peek_data_or_sleep(pipe)
    process_data(data)

This way you won't get thousands of workers when incoming traffic is high.

So the key point is how to send data to worker, usually a pipe or message queue.

I've used Celery with RabbitMQ as message queue. Send data to Celery from Django server, and Celery call your function from another process.

Here is an example assuming you are using asyncio for WebSockets:

import asyncio
from time import sleep

async def web_socket(queue: asyncio.Queue):
    for i in range(5):
        await asyncio.sleep(1.0)
        await queue.put(f"Here is message n°{i}!")
    await queue.put(None)

def expensive_work(message: str):
    sleep(0.5)
    print(message)

async def worker(queue: asyncio.Queue):
    while True:
        message = await queue.get()
        if message is None: break
        await asyncio.to_thread(expensive_work, message)

async def main():
    queue = asyncio.Queue()
    
    await asyncio.gather(
        web_socket(queue),
        worker(queue)
    )

if __name__ == "__main__":
    asyncio.run(main())

The web_socket() function simulates a websocket listener which receives messages. For each received message, it put it in a queue that will be shared with another task running concurrently and processing the message.

The expensive_work() function simulates the processing task to apply to each message.

The worker() function will be running concurrently to the websocket listener. It reads values from the shared queue and process them. If the processing is really expensive (for instance a CPU-bound task) consider running it in a ProcessPoolExecutor (see here how to do that) to avoid blocking the event loop.

Finally, the main() function creates the shared queue, launches the two tasks concurrently with asyncio.gather() and then awaits the completion of both tasks.


If you are using threads and blocking IO, the solution is essentially similar but using threading Threads and queue.Queue. Beware not to mix multithreading and asyncio concurrency, or search on how to do it properly.

Related