diff options
| author | Andy Wingo <wingo@pobox.com> | 2017-01-18 04:27:33 +0100 |
|---|---|---|
| committer | Andy Wingo <wingo@pobox.com> | 2017-01-18 04:32:01 +0100 |
| commit | 4682b331e2ddd4f95740ff319cbdbecc563f2b27 (patch) | |
| tree | dd4d7d10ff7d8b655818190e348f927570a40149 | |
| parent | More speedup tests. (diff) | |
| download | guile-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.scm | 57 | ||||
| -rw-r--r-- | fibers/internal.scm | 82 |
2 files changed, 76 insertions, 63 deletions
@@ -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}." #; |
