Implement server with queue of tasks in python

Viewed 98

I am trying to implement HTTP-server with queue of tasks. See architecture below:

enter image description here

My code

import http.server
import socketserver
import os
import threading
import time

from queue import Queue


PORT = 8004


def run_parser(thread_id, person_id):
    print(f"{threading.current_thread().name} START {thread_id} {person_id}")
    time.sleep(15)
    print(f"{threading.current_thread().name} STOP {thread_id} {person_id}")
    


class MyHttpRequestHandler(http.server.SimpleHTTPRequestHandler):
    def __init__(self, *args, directory=None, **kwargs):
        self.person_ids = Queue()
        self.active_threads = []
        self.THREADS_NUM = 10
        self.all_threads = [] # Threading object instances
        
        super().__init__(*args, directory=None, **kwargs)

    def do_GET(self):
        person_ids = self.path[1:].split(',')
        print("Handle request", person_ids)
        for person_id in person_ids:
            self.process_person(person_id)
    
    def get_free_thread_id(self):
        if len(self.active_threads) == 0:
            return 0
        for i in range(0, self.THREADS_NUM):
            if i not in self.active_threads:
                return i
        return None # no avaliable threads
        
    def process_person(self, person_id):
        """Run thread with person parsing task"""
        thread_id = self.get_free_thread_id()
        print(f"process_person: thread_id={thread_id}, person_id={person_id}")
        if thread_id is not None:
            self.active_threads.append(thread_id)
            print(f"{threading.current_thread().name} QQQ_1", self.active_threads)
            cur_thread = threading.Thread(
                name=f"THREAD_{thread_id}",
                target=run_parser, 
                args=((thread_id, person_id)), 
                daemon=False
            )
            cur_thread.start()
            self.all_threads.append(cur_thread)
        else:
            print(f"Put person_id={person_id} to queue")
            self.person_ids.put(person_id)
            

        

# Create an object of the above class
handler_object = MyHttpRequestHandler


my_server = socketserver.TCPServer(("", PORT), handler_object)

# Star the server
my_server.serve_forever()

Problem For the test

  1. http://127.0.0.1:8004/qwert2,qwert3
  2. After 1 second request http://127.0.0.1:8004/qwert2,qwert3 again

I have the following output:

Handle request ['qwert2', 'qwert3']
process_person: thread_id=0, person_id=qwert2
MainThread QQQ_1 [0]
THREAD_0 START 0 qwert2
process_person: thread_id=1, person_id=qwert3
MainThread QQQ_1 [0, 1]
THREAD_1 START 1 qwert3
Handle request ['qwert2', 'qwert3']
process_person: thread_id=0, person_id=qwert2
MainThread QQQ_1 [0]
THREAD_0 START 0 qwert2
process_person: thread_id=1, person_id=qwert3
MainThread QQQ_1 [0, 1]
THREAD_1 START 1 qwert3
THREAD_0 STOP 0 qwert2
THREAD_1 STOP 1 qwert3
THREAD_0 STOP 0 qwert2
THREAD_1 STOP 1 qwert3

It seems that each do_Get request has its own queue object. But I want ONE queue for all threads and function call

I need output like this (see MainThread QQQ_1 [0, 1, 2, 3]) How can I achieve the output below?

Handle request ['qwert2', 'qwert3']
    process_person: thread_id=0, person_id=qwert2
    MainThread QQQ_1 [0]
    THREAD_0 START 0 qwert2
    process_person: thread_id=1, person_id=qwert3
    MainThread QQQ_1 [0, 1]
    THREAD_1 START 1 qwert3
    Handle request ['qwert2', 'qwert3']
    process_person: thread_id=0, person_id=qwert2
    MainThread QQQ_1 [0, 1, 2]
    THREAD_0 START 0 qwert2
    process_person: thread_id=1, person_id=qwert3
    MainThread QQQ_1 [0, 1, 2, 3]
...
0 Answers
Related