implementing threading to consume queues from rabbitmq

Viewed 9

This is an update to my previous question. I realized that I should have added some code to explain my issue further. I am currently trying to implement threading to queues being consumed from a rabbitmq exchange. As I am new to rabbitmq and threading, I am finding it difficult to amalgamate both and apply them. I was wondering if anyone could provide any templates I could apply to begin with.

I am coding in visual studio platform, where a simulator is being used to generate data, emulating the producer(a smart device).

The first part of the code assigns a few variables relevant to the project. I have imported the required libraries and a few extra scripts I have made myself. The import inthe second line are scripts that assist with communicating with the smart device.

import pika, sys, os
import SockAlertMessage_pb2, SockDataProcessedMessage_pb2, SockDataRawMessage_pb2, SockDataSessionEndMessage_pb2, SockDataSessionStartMessage_pb2, SockMessage_pb2
import numpy as np
import scipy 
import ampd
import python_file_3
import heartpy
import time 
import threading 

from  scipy.signal import detrend
from python_file_3 import filter_signal, get_hrv, get_rmssd, get_std, heart_rate
from ampd import find_peaks_originalode here

The second part of the code

sock_data_session_start_queue = 'sock_data_session_start_queue'
sock_data_session_end_queue= 'sock_data_session_end_queue'
sock_data_raw_queue = 'sock_data_raw_queue'
tx_queue = 'tbd'


# establish connection with rabbitmq server
connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# create rx message queues
channel.queue_declare(queue=sock_data_session_start_queue, durable = True)
channel.queue_declare(queue=sock_data_session_end_queue, durable = True)

# create tx message queue
channel.queue_declare(queue=tx_queue, durable = True)e here


def sock_data_session_start_callback(ch, method, properties, body):

   message = SockDataSessionStartMessage_pb2.SockDataSessionStartMessage()
   message.ParseFromString(body)

   new_dict['1'] = message
   print(new_dict)

   # todo: create thread w/ state
   print(" [x] Start session %r" % message)
    
   # send message

    
def sock_data_session_end_callback(ch, method, properties, body):

   message = SockDataSessionEndMessage_pb2.SockDataSessionEndMessage()
   message.ParseFromString(body)
   # todo: destroy thread w/ state
   print(" [x] End session %r" % message)

def sock_data_raw_callback(ch, method, properties, body):

   message = SockDataRawMessage_pb2.SockDataRawMessage()
   message.ParseFromString(body)
   print(message)

   # todo: destroy thread w/ state
   print(" [x] Sock data raw %r" % message)



if __name__ == '__main__':

    try:
       channel.basic_consume(queue=sock_data_session_start_queue, auto_ack=True, on_message_callback=sock_data_session_start_callback)
       channel.basic_consume(queue=sock_data_session_end_queue, auto_ack=True, on_message_callback=sock_data_session_end_callback)
       channel.basic_consume(queue=sock_data_raw_queue, auto_ack=True, on_message_callback=sock_data_raw_callback)
       print(' [*] Waiting for messages. To exit press CTRL+C')
       channel.start_consuming()

except KeyboardInterrupt:
    print('Interrupted')
    # close connection
    connection.close()
    try:
        sys.exit(0)
    except SystemExit:
        os._exit(0)

The start and end data callbacks refer to data sessions, where data session connection is acknowledged and started and then later ended. I believe the raw data callback is where I may need to implement my threads, where data will be processed and then sent back to another queue. The challenge is to make each data session a thread, and then process that.

0 Answers
Related