summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorAndy Wingo <wingo@pobox.com>2017-01-18 04:27:33 +0100
committerAndy Wingo <wingo@pobox.com>2017-01-18 04:32:01 +0100
commit4682b331e2ddd4f95740ff319cbdbecc563f2b27 (patch)
treedd4d7d10ff7d8b655818190e348f927570a40149
parentMore speedup tests. (diff)
downloadguile-fibers-4682b331e2ddd4f95740ff319cbdbecc563f2b27.tar.gz
Randomized round-robin work sharing/stealing
* fibers.scm (start-auxiliary-threads, stop-auxiliary-threads): Adapt to scheduler-remote-peers change. (spawn-fiber): Adapt to use choose-parallel-scheduler. * fibers/internal.scm (<scheduler>): Add choose-parallel-scheduler field. (shuffle, make-selector): New helpers. (make-scheduler): Adapt to initialize choose-parallel-scheduler field. (choose-parallel-scheduler): New public function. (run-scheduler): Use fiber-stealer.
-rw-r--r--fibers.scm57
-rw-r--r--fibers/internal.scm82
2 files changed, 76 insertions, 63 deletions
diff --git a/fibers.scm b/fibers.scm
index faeddaf..1e8d0ec 100644
--- a/fibers.scm
+++ b/fibers.scm
@@ -63,27 +63,20 @@
(run-scheduler scheduler finished?)))))))
(define (start-auxiliary-threads scheduler hz finished? affinities)
- (let ((scheds (scheduler-remote-peers scheduler)))
- (let lp ((i 0) (affinities affinities))
- (when (< i (vector-length scheds))
- (match affinities
- ((affinity . affinities)
- (let ((remote (vector-ref scheds i)))
- (call-with-new-thread
- (lambda ()
- (%run-fibers remote hz finished? affinity)))
- (lp (1+ i) affinities))))))))
+ (for-each (lambda (sched affinity)
+ (call-with-new-thread
+ (lambda ()
+ (%run-fibers sched hz finished? affinity))))
+ (scheduler-remote-peers scheduler) affinities))
(define (stop-auxiliary-threads scheduler)
- (let ((scheds (scheduler-remote-peers scheduler)))
- (let lp ((i 0))
- (when (< i (vector-length scheds))
- (let* ((remote (vector-ref scheds i))
- (thread (scheduler-kernel-thread remote)))
- (when thread
- (cancel-thread thread)
- (join-thread thread))
- (lp (1+ i)))))))
+ (for-each
+ (lambda (scheduler)
+ (let ((thread (scheduler-kernel-thread scheduler)))
+ (when thread
+ (cancel-thread thread)
+ (join-thread thread))))
+ (scheduler-remote-peers scheduler)))
(define (compute-affinities group-affinity parallelism)
(define (each-thread-has-group-affinity)
@@ -136,29 +129,23 @@
(apply values (atomic-box-ref ret))))))
(define* (spawn-fiber thunk #:optional sched #:key parallel?)
- (define (choose-sched sched)
- (let* ((remote (scheduler-remote-peers sched))
- (count (vector-length remote))
- (idx (random (1+ count))))
- (if (= count idx)
- sched
- (vector-ref remote idx))))
- (define (spawn sched thunk)
- (create-fiber (if parallel? (choose-sched sched) sched)
- thunk))
(cond
(sched
;; When a scheduler is passed explicitly, it could be there is no
;; current fiber; in that case the dynamic state probably doesn't
;; have the right right current-read-waiter /
;; current-write-waiter, so wrap the thunk.
- (spawn sched
- (lambda ()
- (current-read-waiter wait-for-readable)
- (current-write-waiter wait-for-writable)
- (thunk))))
+ (create-fiber sched
+ (lambda ()
+ (current-read-waiter wait-for-readable)
+ (current-write-waiter wait-for-writable)
+ (thunk))))
((current-fiber)
=> (lambda (fiber)
- (spawn (fiber-scheduler fiber) thunk)))
+ (let ((sched (fiber-scheduler fiber)))
+ (create-fiber (if parallel?
+ (choose-parallel-scheduler sched)
+ sched)
+ thunk))))
(else
(error "No scheduler current; call within run-fibers instead"))))
diff --git a/fibers/internal.scm b/fibers/internal.scm
index 99727f6..332b4b5 100644
--- a/fibers/internal.scm
+++ b/fibers/internal.scm
@@ -34,6 +34,7 @@
scheduler-name
(scheduler-kernel-thread/public . scheduler-kernel-thread)
scheduler-remote-peers
+ choose-parallel-scheduler
run-scheduler
destroy-scheduler
@@ -80,7 +81,8 @@ name is known."
(define-record-type <scheduler>
(%make-scheduler name epfd active-fd-count prompt-tag
next-runqueue current-runqueue
- sources timers kernel-thread remote-peers)
+ sources timers kernel-thread
+ remote-peers choose-parallel-scheduler)
scheduler?
(name scheduler-name set-scheduler-name!)
(epfd scheduler-epfd)
@@ -96,8 +98,11 @@ name is known."
(timers scheduler-timers set-scheduler-timers!)
;; atomic parameter of thread
(kernel-thread scheduler-kernel-thread)
- ;; vector of sched
- (remote-peers scheduler-remote-peers set-scheduler-remote-peers!))
+ ;; list of sched
+ (remote-peers scheduler-remote-peers set-scheduler-remote-peers!)
+ ;; () -> sched
+ (choose-parallel-scheduler scheduler-choose-parallel-scheduler
+ set-scheduler-choose-parallel-scheduler!))
(define-record-type <fiber>
(make-fiber scheduler continuation)
@@ -119,6 +124,22 @@ name is known."
(unless (eq? prev init)
(error "owned by other thread" prev))))))))
+(define (shuffle l)
+ (map cdr (sort (map (lambda (x) (cons (random 1.0) x)) l)
+ (lambda (a b) (< (car a) (car b))))))
+
+(define (make-selector items)
+ (let ((items (list->vector (shuffle items))))
+ (match (vector-length items)
+ (0 (lambda () #f))
+ (1 (let ((item (vector-ref items 0))) (lambda () item)))
+ (n (let ((idx 0))
+ (lambda ()
+ (let ((item (vector-ref items idx)))
+ (set! idx (let ((idx (1+ idx)))
+ (if (= idx (vector-length items)) 0 idx)))
+ item)))))))
+
(define* (make-scheduler #:key parallelism
(prompt-tag (make-prompt-tag "fibers")))
"Make a new scheduler in which to run fibers."
@@ -130,24 +151,25 @@ name is known."
(timers (make-psq (match-lambda*
(((t1 . c1) (t2 . c2)) (< t1 t2)))
<))
- (kernel-thread (make-atomic-parameter #f))
- (remote-peers (if parallelism
- (list->vector
- (map (lambda (_)
- (make-scheduler #:prompt-tag prompt-tag))
- (iota (1- parallelism))))
- #())))
+ (kernel-thread (make-atomic-parameter #f)))
(let ((sched (%make-scheduler #f epfd active-fd-count prompt-tag
next-runqueue current-runqueue
sources timers kernel-thread
- remote-peers)))
+ #f #f)))
(set-scheduler-name! sched (nameset-add! schedulers-nameset sched))
- (let lp ((i 0))
- (when (< i (vector-length remote-peers))
- (let ((peers (vector-copy remote-peers)))
- (vector-set! peers i sched)
- (set-scheduler-remote-peers! (vector-ref remote-peers i) peers))
- (lp (1+ i))))
+ (let ((all-scheds
+ (cons sched
+ (if parallelism
+ (map (lambda (_)
+ (make-scheduler #:prompt-tag prompt-tag))
+ (iota (1- parallelism)))
+ '()))))
+ (for-each
+ (lambda (sched)
+ (let ((choose! (make-selector all-scheds)))
+ (set-scheduler-remote-peers! sched (delq sched all-scheds))
+ (set-scheduler-choose-parallel-scheduler! sched choose!)))
+ all-scheds))
sched)))
(define-syntax-rule (with-scheduler scheduler body ...)
@@ -168,6 +190,9 @@ thread."
@code{#f} if @var{sched} is not running."
((scheduler-kernel-thread sched)))
+(define (choose-parallel-scheduler sched)
+ ((scheduler-choose-parallel-scheduler sched)))
+
(define (make-source events expiry fiber) (vector events expiry fiber))
(define (source-events s) (vector-ref s 0))
(define (source-expiry s) (vector-ref s 1))
@@ -275,6 +300,16 @@ thread."
seed))
(run-timers sched))
+(define (fiber-stealer sched)
+ "Steal some work from a random scheduler in the vector
+@var{schedulers}. Return a fiber, or @code{#f} if no work could be
+stolen."
+ (let ((selector (make-selector (scheduler-remote-peers sched))))
+ (lambda ()
+ (let ((peer (selector)))
+ (and peer
+ (stack-pop! (scheduler-current-runqueue peer) #f))))))
+
(define* (run-scheduler sched finished?)
"Run @var{sched} until there are no more fibers ready to run, no
file descriptors being waited on, and no more timers pending to run.
@@ -282,7 +317,7 @@ Return zero values."
(let ((tag (scheduler-prompt-tag sched))
(next (scheduler-next-runqueue sched))
(cur (scheduler-current-runqueue sched))
- (peers (scheduler-remote-peers sched)))
+ (steal-fiber! (fiber-stealer sched)))
(define (run-fiber fiber)
(call-with-prompt tag
(lambda ()
@@ -304,7 +339,7 @@ Return zero values."
;; little bit of work from a remote scheduler if we
;; can. Run it directly instead of pushing onto a
;; queue to avoid double stealing.
- (match (steal-fiber! peers)
+ (match (steal-fiber!)
(#f
(unless (scheduler-finished? sched finished?)
(next-turn)))
@@ -318,15 +353,6 @@ Return zero values."
(run-fiber fiber)
(next-fiber)))))))
-(define (steal-fiber! schedulers)
- "Steal some work from a random scheduler in the vector
-@var{schedulers}. Return a fiber, or @code{#f} if no work could be
-stolen."
- (let ((len (vector-length schedulers)))
- (and (> len 0)
- (let ((sched (vector-ref schedulers (random len))))
- (stack-pop! (scheduler-current-runqueue sched) #f)))))
-
(define (destroy-scheduler sched)
"Release any resources associated with @var{sched}."
#;