3 \P ã @ sb d Z dZddlZddlZddlmZ ddlZddlmZ ddlZddlm Z ddl mZ ddlZddl Z ddlmZ ddlZddlZe jƒ Zd ad d„ ZdZG d d„ deƒZG dd„ dƒZdd„ ZG dd„ deƒZG dd„ deƒZG dd„ deƒZdd„ Zdd„ Z dd„ Z!dd „ Z"d!d"„ Z#d a$da%d#d$„ Z&d%d&„ Z'G d'd(„ d(e(ƒZ)G d)d*„ d*ej*ƒZ+ej,eƒ dS )+a* Implements ProcessPoolExecutor. The follow diagram and text describe the data-flow through the system: |======================= In-process =====================|== Out-of-process ==| +----------+ +----------+ +--------+ +-----------+ +---------+ | | => | Work Ids | => | | => | Call Q | => | | | | +----------+ | | +-----------+ | | | | | ... | | | | ... | | | | | | 6 | | | | 5, call() | | | | | | 7 | | | | ... | | | | Process | | ... | | Local | +-----------+ | Process | | Pool | +----------+ | Worker | | #1..n | | Executor | | Thread | | | | | +----------- + | | +-----------+ | | | | <=> | Work Items | <=> | | <= | Result Q | <= | | | | +------------+ | | +-----------+ | | | | | 6: call() | | | | ... | | | | | | future | | | | 4, result | | | | | | ... | | | | 3, except | | | +----------+ +------------+ +--------+ +-----------+ +---------+ Executor.submit() called: - creates a uniquely numbered _WorkItem and adds it to the "Work Items" dict - adds the id of the _WorkItem to the "Work Ids" queue Local worker thread: - reads work ids from the "Work Ids" queue and looks up the corresponding WorkItem from the "Work Items" dict: if the work item has been cancelled then it is simply removed from the dict, otherwise it is repackaged as a _CallItem and put in the "Call Q". New _CallItems are put in the "Call Q" until "Call Q" is full. NOTE: the size of the "Call Q" is kept small because calls placed in the "Call Q" can no longer be cancelled with Future.cancel(). - reads _ResultItems from "Result Q", updates the future stored in the "Work Items" dict and deletes the dict entry Process #1..n: - reads _CallItems from "Call Q", executes the calls, and puts the resulting _ResultItems in "Result Q" z"Brian Quinlan (brian@sweetapp.com)é N)Ú_base)ÚFull)ÚSimpleQueue)Úwait)ÚpartialFc C sJ da ttjƒ ƒ} x| D ]\}}|jd ƒ qW x| D ]\}}|jƒ q2W d S )NT)Ú _shutdownÚlistÚ_threads_queuesÚitemsÚputÚjoin)r ÚtÚq© r ú2/usr/lib64/python3.6/concurrent/futures/process.pyÚ_python_exitO s r é c @ s e Zd Zdd„ Zdd„ ZdS )Ú_RemoteTracebackc C s || _ d S )N)Útb)Úselfr r r r Ú__init__a s z_RemoteTraceback.__init__c C s | j S )N)r )r r r r Ú__str__c s z_RemoteTraceback.__str__N)Ú__name__Ú __module__Ú__qualname__r r r r r r r ` s r c @ s e Zd Zdd„ Zdd„ ZdS )Ú_ExceptionWithTracebackc C s0 t jt|ƒ||ƒ}dj|ƒ}|| _d| | _d S )NÚ z """ %s""")Ú tracebackÚformat_exceptionÚtyper Úexcr )r r r r r r r g s z _ExceptionWithTraceback.__init__c C s t | j| jffS )N)Ú_rebuild_excr r )r r r r Ú __reduce__l s z"_ExceptionWithTraceback.__reduce__N)r r r r r" r r r r r f s r c C s t |ƒ| _| S )N)r Ú __cause__)r r r r r r! o s r! c @ s e Zd Zdd„ ZdS )Ú _WorkItemc C s || _ || _|| _|| _d S )N)ÚfutureÚfnÚargsÚkwargs)r r% r&