【问题标题】:Python Multithreading Producer Consumer PatternPython 多线程生产者消费者模式
【发布时间】:2019-02-02 02:01:43
【问题描述】:

我仍在学习如何编码,这是我对多线程的第一次尝试。 我读过一堆多线程文章。我认为这些很有帮助:

有很多事情要考虑。特别是对于初学者。 不幸的是,当我尝试将这些信息付诸实践时,我的代码并不能正常工作。

此代码背后的想法是读取包含逗号分隔数字行的simplified.txt。例如:0.275,0.28,0.275,0.275,36078。 生产者线程读取每一行并从行尾删除换行符。然后将行中的每个数字拆分并分配一个变量。 然后将Variable1 放入队列中。 消费者线程将获取队列中的项目,将其平方,然后将条目添加到日志文件中。

我使用的代码来自this template。这是我到目前为止的代码:

import threading
import queue
import time
import logging
import random
import sys

read_file = 'C:/temp/temp1/simplified.txt'
log1 = open('C:/temp/temp1/simplified_log1.txt', "a+")

logging.basicConfig(level=logging.DEBUG, format='(%(threadName)-9s) %(message)s',)

BUF_SIZE = 10
q = queue.Queue(BUF_SIZE)

class ProducerThread(threading.Thread):
    def __init__(self, name, read_file):
        super(ProducerThread,self).__init__()
        self.name = name
        self.read_file = read_file

    def run(self, read_file):
        while True:
            if not q.full():
                with open(read_file, 'r') as f:
                    for line in f:
                        stripped = line.strip('\n\r')
                        value1,value2,value3,value4,value5,value6,value7 = stripped.split(',')
                        q.put(value1)
                        logging.debug('Putting ' + str(value1) + ' : ' + str(q.qsize()) + ' items in queue')
                        time.sleep(random.random())
        return

class ConsumerThread(threading.Thread):
    def __init__(self, name, value1, log1):
        super(ConsumerThread,self).__init__()
        self.name = name
        self.value1 = value1
        self.log1 = log1
        return

    def run(self):
        while True:
            if not q.empty():
                value1 = q.get()
                sqr_value1 = value1 * value1
                log1.write("The square of " + str(value1) + " is " + str(sqr_value1))
                logging.debug('Getting ' + str(value1) + ' : ' + str(q.qsize()) + ' items in queue')
                time.sleep(random.random())
        return

if __name__ == '__main__':
    
    p = ProducerThread(name='producer')
    c = ConsumerThread(name='consumer')

    p.start()
    time.sleep(2)
    c.start()
    time.sleep(2)

当我运行代码时,我得到了这个错误:

Traceback (most recent call last):
  File "c:/Scripta/A_Simplified_Producer_Consumer_Queue_v0.1.py", line 60, in <module>
    p = ProducerThread(name='producer')
TypeError: __init__() missing 1 required positional argument: 'read_file'

我不知道我还需要在哪里添加read_file。 任何帮助将不胜感激。提前致谢。

【问题讨论】:

    标签: python-3.x multithreading


    【解决方案1】:

    您的 ProducerThread 类需要 2 个参数(name 和 read_file)作为其 __init__ 方法中定义的构造函数的参数,当您在 main 中创建实例时,您只提供第一个这样的参数堵塞。你的第二堂课也有同样的问题。

    您应该在创建实例时将read_file 提供给构造函数,或者将其从构造函数签名中删除,因为您似乎并没有使用它(您使用传递给run 函数的read_file,但是我不认为这是正确的)。似乎您正试图从 Thread 超类中覆盖该方法,但我怀疑它采用了这样的参数。

    【讨论】:

      【解决方案2】:

      感谢 userSeventeen 让我走上了正确的道路。 我认为为了使用外部变量,我需要将它们放在 init 方法中,然后再放入 run 方法中。您已经澄清我只需要在运行方法中使用变量。 这是工作代码。我不得不删除 while true: 语句,因为我不希望代码永远运行。

      import threading
      import queue
      import time
      import logging
      import random
      import sys
      import os
      
      
      read_file = 'C:/temp/temp1/simplified.txt'
      log1 = open('C:/temp/temp1/simplified_log1.txt', "a+")
      
      logging.basicConfig(level=logging.DEBUG, format='(%(threadName)-9s) %(message)s',)
      
      BUF_SIZE = 10
      q = queue.Queue(BUF_SIZE)
      
      class ProducerThread(threading.Thread):
          def __init__(self, name):
              super(ProducerThread,self).__init__()
              self.name = name
      
          def run(self):
              with open(read_file, 'r') as f:
                  for line in f:
                      stripped = line.strip('\n\r')
                      value1,value2,value3,value4,value5 = stripped.split(',')
                      float_value1 = float(value1)
                      if not q.full():
                          q.put(float_value1)
                          logging.debug('Putting ' + str(float_value1) + ' : ' + str(q.qsize()) + ' items in queue')
                          time.sleep(random.random())
              return
      
      class ConsumerThread(threading.Thread):
          def __init__(self, name):
              super(ConsumerThread,self).__init__()
              self.name = name
              return
      
          def run(self):
              while not q.empty():
                  float_value1 = q.get()
                  sqr_value1 = float_value1 * float_value1
                  log1.write("The square of " + str(float_value1) + " is " + str(sqr_value1))
                  logging.debug('Getting ' + str(float_value1) + ' : ' + str(q.qsize()) + ' items in queue')
                  time.sleep(random.random())
              return
      
      if __name__ == '__main__':
      
          p = ProducerThread(name='producer')
          c = ConsumerThread(name='consumer')
      
          p.start()
          time.sleep(2)
          c.start()
          time.sleep(2)
      

      【讨论】:

        猜你喜欢
        • 2017-02-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多