Viewing File: /opt/alt/python35/lib64/python3.5/multiprocessing/__pycache__/pool.cpython-35.opt-1.pyc



Yfd@sddgZddlZddlZddlZddlZddlZddlZddlZddlm	Z	ddlm
Z
mZdZdZ
dZejZdd	Zd
dZGdd
d
eZGdddZddZGdddeZdfddddZddZGdddeZGdddeZeZGdddeZGdddeZGd d!d!eZ Gd"ddeZ!dS)#Pool
ThreadPoolN)util)get_contextTimeoutErrorcCstt|S)N)listmap)argsr9/opt/alt/python35/lib64/python3.5/multiprocessing/pool.pymapstar+srcCsttj|d|dS)Nrr)r		itertoolsstarmap)rrrr
starmapstar.src@s(eZdZddZddZdS)RemoteTracebackcCs
||_dS)N)tb)selfrrrr
__init__6szRemoteTraceback.__init__cCs|jS)N)r)rrrr
__str__8szRemoteTraceback.__str__N)__name__
__module____qualname__rrrrrr
r5src@s(eZdZddZddZdS)ExceptionWithTracebackcCsDtjt|||}dj|}||_d||_dS)Nz

"""
%s""")	tracebackformat_exceptiontypejoinexcr)rr rrrr
r<s	zExceptionWithTraceback.__init__cCst|j|jffS)N)rebuild_excr r)rrrr

__reduce__Asz!ExceptionWithTraceback.__reduce__N)rrrrr"rrrr
r;srcCst||_|S)N)r	__cause__)r rrrr
r!Dsr!cs@eZdZdZfddZddZddZS)MaybeEncodingErrorzVWraps possible unpickleable errors, so they can be
    safely sent through the socket.csAt||_t||_tt|j|j|jdS)N)reprr valuesuperr$r)rr r&)	__class__rr
rPszMaybeEncodingError.__init__cCsd|j|jfS)Nz(Error sending result: '%s'. Reason: '%s')r&r )rrrr
rUs	zMaybeEncodingError.__str__cCsd|jj|fS)Nz<%s: %s>)r(r)rrrr
__repr__YszMaybeEncodingError.__repr__)rrr__doc__rrr)rr)r(r
r$Lsr$Fc'Cs|j}|j}t|dr;|jj|jj|dk	rQ||d}x|dksx|r||kry
|}	Wn&ttfk
rtj	dPYnX|	dkrtj	dP|	\}
}}}
}yd||
|f}WnUt
k
rM}z5|r/|tk	r/t||j
}d|f}WYdd}~XnXy||
||fWnbt
k
r}zBt||d}tj	d|||
|d|ffWYdd}~XnXd}	}
}}}
}|d7}qZWtj	d	|dS)
N_writerrz)worker got EOFError or OSError -- exitingzworker got sentinel -- exitingTFrz0Possible encoding error while sending result: %szworker exiting after %d tasks)putgethasattrr+close_readerEOFErrorOSErrorrdebug	Exception_helper_reraises_exceptionr
__traceback__r$)inqueueoutqueueinitializerinitargsZmaxtasksZwrap_exceptionr,r-Z	completedtaskjobifuncrkwdsresultewrappedrrr
worker]sD		


!


	,rCcCs
|dS)z@Pickle-able helper function for use by _guarded_task_generation.Nr)Zexrrr
r5sr5c@seZdZdZdZddZddfddddZdd	Zd
dZdd
Z	ddZ
fiddZdddZdddZ
dddddZddZdddZdddZfidddd Zdddd!d"Zdddd#d$Zed%d&Zed'd(Zed)d*Zed+d,Zd-d.Zd/d0Zd1d2Zd3d4Zed5d6Zed7d8Z d9d:Z!d;d<Z"dS)=rzS
    Class which supports an async version of applying functions to arguments.
    TcOs|jj||S)N)_ctxProcess)rrr?rrr
rEszPool.ProcessNcCs#|pt|_|jtj|_i|_t|_||_	||_
||_|dkrvtj
psd}|dkrtd|dk	rt|rtd||_g|_|jtjdtjd|f|_d|j_t|j_|jjtjdtjd|j|j|j|j|jf|_d|j_t|j_|jjtjdtjd|j|j |jf|_!d|j!_t|j!_|j!jt"j#||j$d|j|j%|j|j|j|j|j!|jfdd|_&dS)	Nrz&Number of processes must be at least 1zinitializer must be a callabletargetrTZexitpriority)'rrD
_setup_queuesqueueQueue
_taskqueue_cacheRUN_state_maxtasksperchild_initializer	_initargsos	cpu_count
ValueErrorcallable	TypeError
_processes_pool_repopulate_pool	threadingZThreadr_handle_workers_worker_handlerdaemonstart
_handle_tasks
_quick_put	_outqueue
_task_handler_handle_results
_quick_get_result_handlerrZFinalize_terminate_pool_inqueue
_terminate)r	processesr9r:Zmaxtasksperchildcontextrrr
rsT
							
		
		
		
z
Pool.__init__cCswd}xjttt|jD]M}|j|}|jdk	r"tjd||jd}|j|=q"W|S)zCleanup after any worker processes which have exited due to reaching
        their specified lifetime.  Returns True if any workers were cleaned up.
        FNzcleaning up worker %dT)reversedrangelenrXexitcoderr3r)rZcleanedr=rCrrr
_join_exited_workerss"

zPool._join_exited_workerscCsxt|jt|jD]}|jdtd|j|j|j|j	|j
|jf}|jj||j
jdd|_
d|_|jtjdqWdS)zBring the number of pool processes up to the specified number,
        for use after reaping workers which have exited.
        rFrrEZ
PoolWorkerTzadded workerN)rlrWrmrXrErCrgrarPrQrO_wrap_exceptionappendnamereplacer]r^rr3)rr=wrrr
rYs#	
zPool._repopulate_poolcCs|jr|jdS)zEClean up any exited workers and start replacements for them.
        N)rorY)rrrr
_maintain_poolszPool._maintain_poolcCsL|jj|_|jj|_|jjj|_|jjj|_	dS)N)
rDZSimpleQueuergrar+sendr`r0recvrd)rrrr
rHszPool._setup_queuescCs|j|||jS)z6
        Equivalent of `func(*args, **kwds)`.
        )apply_asyncr-)rr>rr?rrr
applysz
Pool.applycCs|j||t|jS)zx
        Apply `func` to each element in `iterable`, collecting the results
        in a list that is returned.
        )
_map_asyncrr-)rr>iterable	chunksizerrr
r
szPool.mapcCs|j||t|jS)z
        Like `map()` method but the elements of the `iterable` are expected to
        be iterables as well and will be unpacked as arguments. Hence
        `func` and (a, b) becomes func(a, b).
        )rzrr-)rr>r{r|rrr
rszPool.starmapcCs|j||t|||S)z=
        Asynchronous version of `starmap()` method.
        )rzr)rr>r{r|callbackerror_callbackrrr

starmap_asyncszPool.starmap_asyncccsy>d}x1t|D]#\}}||||fifVqWWn@tk
r}z ||dt|fifVWYdd}~XnXdS)zProvides a generator of tasks for imap and imap_unordered with
        appropriate handling for iterables which throw exceptions during
        iteration.rN)	enumerater4r5)rZ
result_jobr>r{r=xrArrr
_guarded_task_generationszPool._guarded_task_generationrcCs|jtkrtd|dkret|j}|jj|j|j|||j	f|St
j|||}t|j}|jj|j|jt||j	fdd|DSdS)zP
        Equivalent of `map()` -- can be MUCH slower than `Pool.map()`.
        zPool not runningrcss"|]}|D]}|Vq
qdS)Nr).0chunkitemrrr
	<genexpr>@szPool.imap.<locals>.<genexpr>N)
rNrMrTIMapIteratorrLrKr,r_job_set_lengthr
_get_tasksr)rr>r{r|r@task_batchesrrr
imap's 	
	
z	Pool.imapcCs|jtkrtd|dkret|j}|jj|j|j|||j	f|St
j|||}t|j}|jj|j|jt||j	fdd|DSdS)zL
        Like `imap()` method but ordering of results is arbitrary.
        zPool not runningrcss"|]}|D]}|Vq
qdS)Nr)rrrrrr
r[sz&Pool.imap_unordered.<locals>.<genexpr>N)
rNrMrTIMapUnorderedIteratorrLrKr,rrrrrr)rr>r{r|r@rrrr
imap_unorderedBs 	
	
zPool.imap_unorderedcCs_|jtkrtdt|j||}|jj|jd|||fgdf|S)z;
        Asynchronous version of `apply()` method.
        zPool not runningrN)rNrMrTApplyResultrLrKr,r)rr>rr?r}r~r@rrr
rx]s
+zPool.apply_asynccCs|j||t|||S)z9
        Asynchronous version of `map()` method.
        )rzr)rr>r{r|r}r~rrr
	map_asynchszPool.map_asyncc
Cs|jtkrtdt|ds6t|}|dkrztt|t|jd\}}|rz|d7}t|dkrd}tj	|||}t
|j|t||d|}	|jj
|j|	j||df|	S)zY
        Helper function to implement map, starmap and their async counterparts.
        zPool not running__len__Nrrr~)rNrMrTr.r	divmodrmrXrr	MapResultrLrKr,rr)
rr>r{Zmapperr|r}r~Zextrarr@rrr
rzps&(
		
zPool._map_asynccCsrtj}xB|jtks6|jrP|jtkrP|jtjdqW|j	j
dtjddS)Ng?zworker handler exiting)
rZcurrent_threadrNrMrL	TERMINATErutimesleeprKr,rr3)poolthreadrrr
r[s*
zPool._handle_workersc
Cstj}x+t|jdD]
\}}d}zx|D]}|jrXtjdPy||Wq;tk
r}	zN|dd\}
}y||
j|d|	fWnt	k
rYnXWYdd}	~	Xq;Xq;W|rtjd|r|dnd}||dwPWdd}}}
XqWtjdyFtjd|j
dtjdx|D]}|dqkWWntk
rtjd	YnXtjd
dS)Nz'task handler found thread._state != RUNrFzdoing set_length()rztask handler got sentinelz/task handler sending sentinel to result handlerz(task handler sending sentinel to workersz/task handler got OSError when sending sentinelsztask handler exitingr)rZriterr-rNrr3r4_setKeyErrorr,r2)
	taskqueuer,r8rcacherZtaskseqZ
set_lengthr;rAr<idxprrr
r_sB
	








zPool._handle_taskscCstj}xy
|}Wn)ttfk
rGtjddSYnX|jr_tjdP|dkrytjdP|\}}}y||j||Wntk
rYnXd}}}qWx|r|jt	kry
|}Wn)ttfk
rtjddSYnX|dkr4tjdq|\}}}y||j||Wntk
roYnXd}}}qWt
|drtjdy2x+tdD]}|jj
sP|qWWnttfk
rYnXtjdt||jdS)	Nz.result handler got EOFError/OSError -- exitingz,result handler found thread._state=TERMINATEzresult handler got sentinelz&result handler ignoring extra sentinelr0z"ensuring that outqueue is not full
z7result handler exiting: len(cache)=%s, thread._state=%s)rZrr2r1rr3rNrrrr.rlr0pollrm)r8r-rrr;r<r=objrrr
rcsZ

		




	


	zPool._handle_resultsccsDt|}x1ttj||}|s1dS||fVqWdS)N)rtuplerislice)r>itsizerrrr
rszPool._get_taskscCstddS)Nz:pool objects cannot be passed between processes or pickled)NotImplementedError)rrrr
r"szPool.__reduce__cCs5tjd|jtkr1t|_t|j_dS)Nzclosing pool)rr3rNrMCLOSEr\)rrrr
r/s
	z
Pool.closecCs0tjdt|_t|j_|jdS)Nzterminating pool)rr3rrNr\rh)rrrr
	terminates
	zPool.terminatecCsVtjd|jj|jj|jjx|jD]}|jq>WdS)Nzjoining pool)rr3r\rrbrerX)rrrrr
rs



z	Pool.joincCsZtjd|jjx9|jrU|jjrU|jjtj	dqWdS)Nz7removing tasks from inqueue until task handler finishedr)
rr3Z_rlockacquireis_aliver0rrwrr)r7task_handlerrrrr
_help_stuff_finish(s



zPool._help_stuff_finishc	
Cstjdt|_t|_tjd|j||t|t|_|jdtjdtj|k	r|j	|rt
|ddrtjdx'|D]}	|	jdkr|	jqWtjdtj|k	r|j	tjdtj|k	r&|j	|rt
|ddrtjd	x8|D]0}	|	j
rStjd
|	j|	j	qSWdS)Nzfinalizing poolz&helping task handler/workers to finishzjoining worker handlerrrzterminating workerszjoining task handlerzjoining result handlerzjoining pool workerszcleaning up worker %d)rr3rrNrrmr,rZrrr.rnrrpid)
clsrr7r8rZworker_handlerrZresult_handlerrrrrr
rf1s6
		
	










zPool._terminate_poolcCs|S)Nr)rrrr
	__enter___szPool.__enter__cCs|jdS)N)r)rexc_typeZexc_valZexc_tbrrr
__exit__bsz
Pool.__exit__)#rrrr*rprErrorYrurHryr
rrrrrrxrrzstaticmethodr[r_rcrr"r/rrrclassmethodrfrrrrrr
rsF	8	

.<			.c@s^eZdZddZddZddZddd	Zdd
dZdd
ZdS)rcCsJtj|_tt|_||_||_||_|||j<dS)N)	rZZEvent_eventnextjob_counterrrL	_callback_error_callback)rrr}r~rrr
rks			zApplyResult.__init__cCs
|jjS)N)rZis_set)rrrr
readysszApplyResult.readycCs|jS)N)_success)rrrr

successfulvszApplyResult.successfulNcCs|jj|dS)N)rwait)rtimeoutrrr
rzszApplyResult.waitcCs<|j||jst|jr/|jS|jdS)N)rrrr_value)rrrrr
r-}s
	zApplyResult.getcCsu|\|_|_|jr4|jr4|j|j|jrW|jrW|j|j|jj|j|j=dS)N)rrrrrsetrLr)rr=rrrr
rs
zApplyResult._set)	rrrrrrrr-rrrrr
ris	rc@s(eZdZddZddZdS)rcCstj|||d|d|_dg||_||_|dkrjd|_|jj||j=n||t	|||_dS)Nr~Tr)
rrrr
_chunksize_number_leftrrrbool)rrr|lengthr}r~rrr
rs			

zMapResult.__init__cCs|\}}|r||j||j|d|j<|jd8_|jdkr|jrn|j|j|j|j=|jjnEd|_||_|j	r|j	|j|j|j=|jjdS)NrrF)
rrrrrLrrrrr)rr=Zsuccess_resultsuccessr@rrr
rs%	
			
zMapResult._setN)rrrrrrrrr
rs
rc@sUeZdZddZddZdddZeZdd	Zd
dZdS)rcCsktjtj|_tt|_||_tj	|_
d|_d|_i|_
|||j<dS)Nr)rZZ	ConditionZLock_condrrrrLcollectionsdeque_items_index_length	_unsorted)rrrrr
rs				zIMapIterator.__init__cCs|S)Nr)rrrr
__iter__szIMapIterator.__iter__NcCs|jy|jj}Wntk
r|j|jkrEt|jj|y|jj}Wn0tk
r|j|jkrttYnXYnXWdQRX|\}}|r|S|dS)N)	rrpopleft
IndexErrorrr
StopIterationrr)rrrrr&rrr
rs"


zIMapIterator.nextc
Cs|j|j|kr|jj||jd7_xJ|j|jkr|jj|j}|jj||jd7_q;W|jjn
||j|<|j|jkr|j|j	=WdQRXdS)Nr)
rrrrqrpopnotifyrrLr)rr=rrrr
rs

zIMapIterator._setc	CsJ|j:||_|j|jkr?|jj|j|j=WdQRXdS)N)rrrrrLr)rrrrr
rs

	
zIMapIterator._set_length)	rrrrrr__next__rrrrrr
rs
rc@seZdZddZdS)rc
Cs`|jP|jj||jd7_|jj|j|jkrU|j|j=WdQRXdS)Nr)rrrqrrrrLr)rr=rrrr
rs

zIMapUnorderedIterator._setN)rrrrrrrr
rsrc@s[eZdZdZeddZddfddZddZed	d
ZdS)rFcOsddlm}|||S)Nr)rE)ZdummyrE)rr?rErrr
rEszThreadPool.ProcessNcCstj||||dS)N)rr)rrir9r:rrr
rszThreadPool.__init__cCs@tj|_tj|_|jj|_|jj|_dS)N)rIrJrgrar,r`r-rd)rrrr
rHszThreadPool._setup_queuesc
CsF|j6|jj|jjdg||jjWdQRXdS)N)Z	not_emptyrIclearextendZ
notify_all)r7rrrrr
rs

zThreadPool._help_stuff_finish)	rrrrprrErrHrrrrr
rs
)"__all__rZrIrrrRrrrrrrrMrrcountrrrr4rrr!r$rCr5objectrrZAsyncResultrrrrrrrr
<module>
s<		*&%@
Back to Directory File Manager