【问题标题】:Using Queues in Lparallel Library (Common Lisp)在 Lparallel 库中使用队列(Common Lisp)
【发布时间】:2019-08-14 13:31:37
【问题描述】:

https://z0ltan.wordpress.com/2016/09/09/basic-concurrency-and-parallelism-in-common-lisp-part-4a-parallelism-using-lparallel-fundamentals/#channels 的 lparallel 库中关于队列的基本讨论说队列“允许在工作线程之间传递消息”。下面的测试使用一个共享队列来协调一个主线程和一个从属线程,其中主线程只是在退出之前等待从属线程完成:

(defun foo (q)
  (sleep 1)
  (lparallel.queue:pop-queue q))  ;q is now empty

(defun test ()
  (setf lparallel:*kernel* (lparallel:make-kernel 1))
  (let ((c (lparallel:make-channel))
        (q (lparallel.queue:make-queue)))
    (lparallel.queue:push-queue 0 q)
    (lparallel:submit-task c #'foo q)
    (loop do (sleep .2)
             (print (lparallel.queue:peek-queue q))
          when (lparallel.queue:queue-empty-p q)
            do (return)))
  (lparallel:end-kernel :wait t))

这按预期产生输出:

* (test)

0
0
0
0
NIL
(#<SB-THREAD:THREAD "lparallel" FINISHED values: NIL {10068F2B03}>)

我的问题是我是否正确或完全使用 lparallel 的队列功能。看起来队列只是使用全局变量来保存线程共享对象的替代品。使用队列的设计优势是什么?为每个提交的任务分配一个队列(假设任务需要通信)通常是一种好习惯吗?感谢您提供更深入的见解。

【问题讨论】:

    标签: multithreading common-lisp message-queue


    【解决方案1】:

    多线程工作是通过管理对 mutable 的并发访问来完成的 共享状态,即你有一个通用数据结构的锁, 每个线程读取或写入它。

    但是建议尽量减少数据的数量 同时访问。队列是一种将工作人员与每个工作人员分离的方法 其他,通过让每个线程管理其本地状态并交换数据 仅通过消息;这是线程安全的,因为访问 队列由locks and condition variables控制。

    您在主线程中所做的是轮询队列何时 是空的;这可能会起作用,但这会适得其反,因为队列 被用作同步机制,但在这里你正在做 自己同步。

    (ql:quickload :lparallel)
    (defpackage :so (:use :cl
                          :lparallel
                          :lparallel.queue
                          :lparallel.kernel-util))
    (in-package :so)
    

    让我们改变foo,让它有两个队列,一个用于传入 请求,一个用于回复。在这里,我们执行一个简单的转换为 正在发送的数据,对于每个输入消息,只有一个 输出消息,但不一定总是如此。

    (defun foo (in out)
      (push-queue (1+ (pop-queue in)) out))
    

    更改test,使控制流仅基于对队列的读/写:

    (defun test ()
      (with-temp-kernel (1)
        (let ((c (make-channel))
              (foo-in (make-queue))
              (foo-out (make-queue)))
          (submit-task c #'foo foo-in foo-out)
          ;; submit data to task (could be blocking)
          (push-queue 0 foo-in)
          ;; wait for message from task (could be blocking too)
          (pop-queue foo-out))))
    

    但是,如果有多个任务正在运行,如何避免在测试中进行轮询?您是否不需要不断检查其中任何一项何时完成,以便您可以将更多工作推入队列?

    您可以使用不同的并发机制,类似于 listenpoll/epoll,您可以在其中监视多个 事件的来源,并在其中一个准备好时做出反应。有像 Go (select) 和 Erlang (receive) 这样的语言 这是很自然的表达方式。在 Lisp 方面,Calispel 库提供了类似的交替机制(pri-altfair-alt)。例如,以下内容来自 Calispel 的测试代码:

    (pri-alt ((? control msg)
              (ecase msg
                (:clean-up (setf cleanup? t))
                (:high-speed (setf slow? nil))
                (:low-speed (setf slow? t))))
             ((? channel msg)
              (declare (type fixnum msg))
              (vector-push-extend msg out))
             ((otherwise :timeout (if cleanup? 0 nil))
              (! reader-results out)
              (! thread-expiration (bt:current-thread))
              (return)))
    

    lparallel 的情况下,没有这样的机制,但你可以只使用队列,只要你用标识符标记你的消息。

    如果您需要在任一任务t1t2 给出结果后立即做出反应,那么让这两个任务写入同一个结果通道:

    (let ((t1 (foo :id 1 :in i1 :out res))
          (t2 (bar :id 2 :in i2 :out res)))
       (destructuring-bind (id message) (pop-queue res)
         (case id
           (1 ...)
           (2 ...))))
    

    如果你需要在t1t2发出结果时同步代码,让它们在不同的通道中写入:

    (let ((t1 (foo :id 1 :in i1 :out o1))
          (t2 (bar :id 2 :in i2 :out o2)))
       (list (pop-queue o1)
             (pop-queue o2)))
    

    【讨论】:

    • 那么,队列大致类似于函数的 lambda 列表?即,一个函数可以将任意 lisp 对象(作为参数)传递给一个从属函数,一个线程可以将任意 lisp 对象(作为“out”队列元素)传递给另一个线程。然后接收线程可以使用“in”队列元素来初始化或以其他方式执行其任务(或阻塞直到元素变得可用)。接收线程也可以通过使用它的“out”队列来请求额外的数据,并在它的“in”队列中接收数据。
    • 最终,接收线程可以通过使用它的“out”队列来报告它的结果(或者在 lparallel 中通过它的通道返回一个值)。这有意义吗?
    • 但是如果有多个任务正在运行,如何避免在 test 中进行轮询?您不需要不断检查其中任何一项何时完成,以便您可以 push-queue 完成更多工作吗?
    • 直截了当的解释!您必须是 CS 教授。我想我现在会坚持使用 lparallel,但很高兴了解其他人(魂器问题)。
    • 谢谢!祝你好运!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-11-10
    • 2011-03-31
    • 2014-11-28
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多