o
    l¨Êhµ(  ã                   @   sL  d dl mZmZmZ d dlZd dlmZ d dlZddlm	Z	m
Z
mZ ddlmZ zd dlZW n ey9   dZY nw edurjd dlmZmZ d dlmZmZ d d	lmZmZmZ d d
lmZ d dlmZ d dlmZ dd„ ZzeZW n e yƒ   G dd„ de!ƒZY nw G dd„ dƒZ"dd„ Z#G dd„ de$ƒZ%dd„ Z&G dd„ de
e	ƒZ'dS )é    )Úprint_functionÚdivisionÚabsolute_importN)Úuuid4é   )ÚAutoBatchingMixinÚParallelBackendBaseÚBatchedCalls)Úparallel_backend)ÚClientÚ_wait)ÚfuncnameÚ
itemgetter)Ú
get_clientÚsecedeÚrejoin)Úthread_state)Úsizeof)Úgenc                 C   s&   zt  | ¡ W dS  ty   Y dS w )NTF)ÚweakrefÚrefÚ	TypeError)Úobj© r   ú>/var/www/html/env/lib/python3.10/site-packages/joblib/_dask.pyÚis_weakrefable   s   
ÿr   c                   @   s   e Zd ZdS )ÚTimeoutErrorN)Ú__name__Ú
__module__Ú__qualname__r   r   r   r   r   %   s    r   c                   @   s8   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ ZdS )Ú_WeakKeyDictionarya«  A variant of weakref.WeakKeyDictionary for unhashable objects.

    This datastructure is used to store futures for broadcasted data objects
    such as large numpy arrays or pandas dataframes that are not hashable and
    therefore cannot be used as keys of traditional python dicts.

    Futhermore using a dict with id(array) as key is not safe because the
    Python is likely to reuse id of recently collected arrays.
    c                 C   s
   i | _ d S ©N©Ú_data©Úselfr   r   r   Ú__init__4   ó   
z_WeakKeyDictionary.__init__c                 C   s(   | j t|ƒ \}}|ƒ |urt|ƒ‚|S r!   )r#   ÚidÚKeyError)r%   r   r   Úvalr   r   r   Ú__getitem__7   s   
z_WeakKeyDictionary.__getitem__c                    sl   t |ƒ‰ zˆjˆ  \}}|ƒ |urt|ƒ‚W n ty,   ‡ ‡fdd„}t ||¡}Y nw ||fˆjˆ < d S )Nc                    s   ˆj ˆ = d S r!   r"   )Ú_©Úkeyr%   r   r   Ú
on_destroyI   ó   z2_WeakKeyDictionary.__setitem__.<locals>.on_destroy)r(   r#   r)   r   r   )r%   r   Úvaluer   r,   r/   r   r-   r   Ú__setitem__>   s   
þúz_WeakKeyDictionary.__setitem__c                 C   s
   t | jƒS r!   )Úlenr#   r$   r   r   r   Ú__len__N   r'   z_WeakKeyDictionary.__len__c                 C   ó   | j  ¡  d S r!   )r#   Úclearr$   r   r   r   r6   Q   s   z_WeakKeyDictionary.clearN)	r   r   r   Ú__doc__r&   r+   r2   r4   r6   r   r   r   r   r    )   s    
r    c                 C   sF   zt | tƒr| jd d } W t| ƒS W t| ƒS  ty"   Y t| ƒS w )Nr   )Ú
isinstancer	   ÚitemsÚ	Exceptionr   )Úxr   r   r   Ú	_funcnameU   s   
üþþr<   c                   @   s$   e Zd Zdd„ Zdd„ Zdd„ ZdS )ÚBatchc                 C   s
   || _ d S r!   )Útasks)r%   r>   r   r   r   r&   _   r'   zBatch.__init__c                    s€   g }t dƒ�0 | jD ]#\}}}‡ fdd„|D ƒ}‡ fdd„| ¡ D ƒ}| ||i |¤Ž¡ q
W d   ƒ |S 1 s9w   Y  |S )NÚdaskc                    s"   g | ]}t |tƒr|ˆ ƒn|‘qS r   ©r8   r   )Ú.0Úa©Údatar   r   Ú
<listcomp>f   s    ÿz"Batch.__call__.<locals>.<listcomp>c                    s(   i | ]\}}|t |tƒr|ˆ ƒn|“qS r   r@   )rA   ÚkÚvrC   r   r   Ú
<dictcomp>h   s    ÿz"Batch.__call__.<locals>.<dictcomp>)r
   r>   r9   Úappend)r%   rD   ÚresultsÚfuncÚargsÚkwargsr   rC   r   Ú__call__b   s   

ÿ
ÿû
ÿùzBatch.__call__c                 C   s   t | jffS r!   )r=   r>   r$   r   r   r   Ú
__reduce__m   r0   zBatch.__reduce__N)r   r   r   r&   rN   rO   r   r   r   r   r=   ^   s    r=   c                   C   s   d S r!   r   r   r   r   r   Ú_joblib_probe_taskq   s   rP   c                   @   s~   e Zd ZdZdZ		ddd„Zdd„ Zd	d
„ Zddd„Zdd„ Z	dd„ Z
dd„ Zdd„ Zddd„Zd dd„Zejdd„ ƒZdS )!ÚDaskDistributedBackendgš™™™™™É?g      ð?Né
   c           	      K   sî   t d u r
d}t|ƒ‚|d u r+|rt||dd�}nztƒ }W n ty*   d}t|ƒ‚w || _|d urBt|ttfƒsBtdt	|ƒj
 ƒ‚|d uret|ƒdkret|ƒ| _| jj|dd�}d	d
„ t||ƒD ƒ| _ng | _i | _tƒ | _|| _|| _d S )Nz{You are trying to use 'dask' as a joblib parallel backend but dask is not installed. Please install dask to fix this error.F)ÚloopÚset_as_defaultz¢To use Joblib with Dask first create a Dask Client

    from dask.distributed import Client
    client = Client()
or
    client = Client('scheduler-address:8786')z&scatter must be a list/tuple, got `%s`r   T)Ú	broadcastc                 S   s   i | ]	\}}t |ƒ|“qS r   )r(   )rA   r;   Úfr   r   r   rH   �   s    z3DaskDistributedBackend.__init__.<locals>.<dictcomp>)ÚdistributedÚ
ValueErrorr   r   Úclientr8   ÚlistÚtupler   Útyper   r3   Ú_scatterÚscatterÚzipÚdata_futuresÚsetÚtask_futuresÚwait_for_workers_timeoutÚsubmit_kwargs)	r%   Úscheduler_hostr^   rY   rS   rc   rd   ÚmsgÚ	scatteredr   r   r   r&   z   s8   ÿ
ù	ÿ

zDaskDistributedBackend.__init__c                 C   s   t dfS )Nr   )rQ   r$   r   r   r   rO   ¥   s   z!DaskDistributedBackend.__reduce__c                 C   s   t | jd�dfS )N)rY   éÿÿÿÿ)rQ   rY   r$   r   r   r   Úget_nested_backend¨   s   z)DaskDistributedBackend.get_nested_backendr   c                 K   s
   |   |¡S r!   )Úeffective_n_jobs)r%   Ún_jobsÚparallelÚbackend_argsr   r   r   Ú	configure«   r'   z DaskDistributedBackend.configurec                 C   s   t ƒ | _d S r!   )r    Úcall_data_futuresr$   r   r   r   Ú
start_call®   r0   z!DaskDistributedBackend.start_callc                 C   r5   r!   )ro   r6   r$   r   r   r   Ú	stop_call±   s   z DaskDistributedBackend.stop_callc              
   C   s„   t | j ¡  ¡ ƒ}|dks| js|S z| j t¡j| jd� W n tj	y8   d 
| jtdd| j ƒ¡}t	|ƒ‚w t | j ¡  ¡ ƒS )Nr   )ÚtimeoutzðDaskDistributedBackend has no worker after {} seconds. Make sure that workers are started and can properly connect to the scheduler and increase the joblib/dask connection timeout with:

parallel_backend('dask', wait_for_workers_timeout={})rR   é   )ÚsumrY   ÚncoresÚvaluesrc   ÚsubmitrP   Úresultr   r   ÚformatÚmax)r%   rk   rj   Ú	error_msgr   r   r   rj   ¶   s    
ÿÿú÷
z'DaskDistributedBackend.effective_n_jobsc                    sŒ   g ‰t ƒ ‰tˆdd ƒ‰ ‡ ‡‡‡fdd„}g }|jD ] \}}}t||ƒƒ}t t| ¡ || ¡ ƒƒƒ}| |||f¡ qˆs@|dfS t|ƒˆfS )Nro   c              	   3   sÆ   � | D ]]}t |ƒ}|ˆv rˆ| V  qˆj |d ¡}|d u rHˆ d urHzˆ | }W n tyG   t|ƒrEt|ƒdkrEˆj |g¡\}|ˆ |< Y nw |d ur]tt	ˆƒƒ}ˆ 
|¡ |ˆ|< |}|V  qd S )Ng     @�@)r(   r`   Úgetr)   r   r   rY   r^   r   r3   rI   )rL   ÚargÚarg_idrV   Úgetter©ro   Úcollected_futuresÚitemgettersr%   r   r   Úmaybe_to_futuresÕ   s.   €
€ø

çz>DaskDistributedBackend._to_func_args.<locals>.maybe_to_futuresr   )	ÚdictÚgetattrr9   rZ   r_   Úkeysrv   rI   r=   )r%   rK   rƒ   r>   rV   rL   rM   r   r€   r   Ú_to_func_argsÍ   s   

ÿz$DaskDistributedBackend._to_func_argsc                    s†   dt |ƒtƒ jf }ˆ |¡\}}ˆjj|g|¢R d|iˆj¤Ž}ˆj |¡ ‡ ‡fdd„}| 	|¡ t
 |¡‰‡fdd„}||_|S )Nz%s-batch-%sr.   c                    s,   |   ¡ }ˆj | ¡ ˆ d urˆ |ƒ d S d S r!   )rx   rb   Úremove)Úfuturerx   )Úcallbackr%   r   r   Úcallback_wrapper  s
   ÿz<DaskDistributedBackend.apply_async.<locals>.callback_wrapperc                      s
   ˆ ƒ   ¡ S r!   )rx   r   )r   r   r   r|     r'   z/DaskDistributedBackend.apply_async.<locals>.get)r<   r   Úhexr‡   rY   rw   rd   rb   ÚaddÚadd_done_callbackr   r   r|   )r%   rK   rŠ   r.   rL   r‰   r‹   r|   r   )rŠ   r   r%   r   Úapply_asyncü   s    

z"DaskDistributedBackend.apply_asyncTc                 C   s   | j  | j¡ | j ¡  dS )z� Tell the client to cancel any task submitted via this instance

        joblib.Parallel will never access those results
        N)rY   Úcancelrb   r6   )r%   Úensure_readyr   r   r   Úabort_everything  s   z'DaskDistributedBackend.abort_everythingc                 c   s0   � t tdƒr	tƒ  dV  t tdƒrtƒ  dS dS )zÙOverride ParallelBackendBase.retrieval_context to avoid deadlocks.

        This removes thread from the worker's thread pool (using 'secede').
        Seceding avoids deadlock in nested parallelism settings.
        Úexecution_stateN)Úhasattrr   r   r   r$   r   r   r   Úretrieval_context  s   €
	

ÿz(DaskDistributedBackend.retrieval_context)NNNNrR   )r   Nr!   )T)r   r   r   ÚMIN_IDEAL_BATCH_DURATIONÚMAX_IDEAL_BATCH_DURATIONr&   rO   ri   rn   rp   rq   rj   r‡   r�   r’   Ú
contextlibÚcontextmanagerr•   r   r   r   r   rQ   v   s"    
ÿ+

/
rQ   )(Ú
__future__r   r   r   r˜   Úuuidr   r   rl   r   r   r	   r
   rW   ÚImportErrorÚdistributed.clientr   r   Údistributed.utilsr   r   r   r   r   Údistributed.workerr   Údistributed.sizeofr   Útornador   r   r   Ú	NameErrorÚOSErrorr    r<   Úobjectr=   rP   rQ   r   r   r   r   Ú<module>   s:    ÿþ,	