403Webshell
Server IP : 54.37.205.81  /  Your IP : 216.73.216.76
Web Server : nginx/1.22.1
System : Linux vps-249481fa 6.1.0-50-cloud-amd64 #1 SMP PREEMPT_DYNAMIC Debian 6.1.176-1 (2026-07-02) x86_64
User : debian ( 1000)
PHP Version : 8.2.32
Disable Function : NONE
MySQL : OFF  |  cURL : ON  |  WGET : ON  |  Perl : ON  |  Python : OFF  |  Sudo : ON  |  Pkexec : OFF
Directory :  /lib/python3/dist-packages/joblib/__pycache__/

Upload File :
current_dir [ Writeable ] document_root [ Writeable ]

 

Command :


[ Back ]     

Current File : /lib/python3/dist-packages/joblib/__pycache__/_dask.cpython-311.pyc
�


I$c�2���ddlmZmZmZddlZddlZddlZddlZddl	m
Z
ddlZddlm
Z
mZmZddlmZ	ddlZddlZn#e$rdZdZYnwxYwe�?e�=ddlmZmZddlmZdd	lmZmZmZmZmZmZdd
l m!Z!	ddl m"Z#n#e$r	ddl$m"Z#YnwxYwd�Z%Gd
�d��Z&d�Z'd�Z(Gd�d��Z)d�Z*Gd�de
e��Z+dS)�)�print_function�division�absolute_importN)�uuid4�)�AutoBatchingMixin�ParallelBackendBase�BatchedCalls)�parallel_backend)�funcname�
itemgetter)�sizeof)�Client�as_completed�
get_client�secede�rejoin�
get_worker)�thread_state)�TimeoutErrorc�R�	tj|��dS#t$rYdSwxYw)NTF)�weakref�ref�	TypeError)�objs �./usr/lib/python3/dist-packages/joblib/_dask.py�is_weakrefabler+s>�����C�����t�������u�u����s��
&�&c�0�eZdZdZd�Zd�Zd�Zd�Zd�ZdS)�_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.

    Furthermore using a dict with id(array) as key is not safe because the
    Python is likely to reuse id of recently collected arrays.
    c��i|_dS�N��_data��selfs r�__init__z_WeakKeyDictionary.__init__>s
����
�
�
�c�v�|jt|��\}}|��|urt|���|Sr!)r#�id�KeyError)r%rr�vals    r�__getitem__z_WeakKeyDictionary.__getitem__As:���:�b��g�g�&���S��3�5�5�����3�-�-���
r'c�����t|���	�j�\}}|��|urt|���n+#t$r��fd�}tj||��}YnwxYw||f�j�<dS)Nc����j�=dSr!r")�_�keyr%s ��r�
on_destroyz2_WeakKeyDictionary.__setitem__.<locals>.on_destroySs����J�s�O�O�Or')r)r#r*rr)r%r�valuerr/r1r0s`     @r�__setitem__z_WeakKeyDictionary.__setitem__Hs�������g�g��	/��Z��_�F�C���s�u�u�C����s�m�m�#� ���	/�	/�	/�
$�
$�
$�
$�
$�
$��+�c�:�.�.�C�C�C�
	/�����u�*��
�3���s�+?�%A'�&A'c�*�t|j��Sr!)�lenr#r$s r�__len__z_WeakKeyDictionary.__len__Xs���4�:���r'c�8�|j���dSr!)r#�clearr$s rr8z_WeakKeyDictionary.clear[s���
�������r'N)	�__name__�
__module__�__qualname__�__doc__r&r,r3r6r8�r'rrr3si��������������%�%�%� �������r'rc��	t|t��r|dd}n#t$rYnwxYwt|��S)Nr)�
isinstance�list�	Exceptionr)�xs r�	_funcnamerC_sS��
��a����	��!��Q��A����
�
�
���
�����A�;�;�s�#&�
3�3c��d�|D��}t|��dkrd}nd}t|��|t|��fS)z8Summarize of list of (func, args, kwargs) function callsc��h|]\}}}|��	Sr=r=)�.0�func�args�kwargss    r�	<setcomp>z&_make_tasks_summary.<locals>.<setcomp>js��9�9�9�/�T�4��D�9�9�9r'rFT)r5rC)�tasks�unique_funcs�mixeds   r�_make_tasks_summaryrNhsO��9�9�5�9�9�9�L�
�<���A���������u�:�:�u�i��.�.�.�.r'c�&�eZdZdZd�Zdd�Zd�ZdS)�Batchz6dask-compatible wrapper that executes a batch of tasksc�J�t|��\|_|_|_dSr!)rN�
_num_tasks�_mixedrC)r%rKs  rr&zBatch.__init__us(��8K��8
�8
�4�����d�n�n�nr'Nc	��g}td��5|D]!\}}}|�||i|�����"	ddd��n#1swxYwY|S)N�dask)r�append)r%rK�resultsrGrHrIs      r�__call__zBatch.__call__|s�����
�f�
%�
%�	6�	6�&+�
6�
6�"��d�F����t�t�T�4�V�4�4�5�5�5�5�
6�	6�	6�	6�	6�	6�	6�	6�	6�	6�	6�	6����	6�	6�	6�	6��s�%A�A�Ac�D�d|j�d|j�d�}|jrd|z}|S)N�	batch_of_r/�_calls�mixed_)rCrRrS)r%�descrs  r�__repr__zBatch.__repr__�s7��D�D�N�D�D�T�_�D�D�D���;�	%��u�$�E��r'r!)r9r:r;r<r&rXr^r=r'rrPrPssL������@�@�
�
�
���������r'rPc��dSr!r=r=r'r�_joblib_probe_taskr`�s���Dr'c���eZdZdZdZdZ		d�fd�	Zd�Zd�Zd	�Z	dd�Z
d�Zd
�Zd�Z
d�Zdd�Zdd�Zejd���Z�xZS)�DaskDistributedBackendg�������?g�?TN�
c���t�����t�d}t|���|�I|rt	||d���}n4	t��}n$#t$r}d}t|��|�d}~wwxYw||_|�@t|ttf��s$tdt|��jz���|�ct|��dkrPt|��|_|j�|d���}	d	�t!||	��D��|_ng|_i|_||_||_t)g|jdd�
��|_i|_i|_dS)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`rT)�	broadcastc�4�i|]\}}t|��|��Sr=)r))rFrB�fs   r�
<dictcomp>z3DaskDistributedBackend.__init__.<locals>.<dictcomp>�s$�� N� N� N�d�a���A���� N� N� Nr')re�with_results�raise_errors)�superr&�distributed�
ValueErrorrr�clientr?r@�tupler�typer9r5�_scatter�scatter�zip�data_futures�wait_for_workers_timeout�
submit_kwargsrre�waiting_futures�_results�
_callbacks)r%�scheduler_hostrtrprerwrx�msg�e�	scattered�	__class__s          �rr&zDaskDistributedBackend.__init__�s����	����������(�C��S�/�/�!��>��

1���T�/4�6�6�6���	1�'�\�\�F�F��!�1�1�1�K�C�%�S�/�/�q�0�����1���������z�'�D�%�=�'I�'I���#�%)�'�]�]�%;�<�=�=�
=���3�w�<�<�!�#3�#3� ��M�M�D�M���+�+�G�t�+�D�D�I� N� N�c�'�9�6M�6M� N� N� N�D����D�M� "�D��(@��%�*���+������	 
� 
� 
�����
�����s�A!�!
B�+A=�=Bc��zK�|jr�|j23d{V��\}}|j�|��}|j�|��}|jdkr|\}}}|�|���e|�|��||����6tj	d���d{V��|j��dSdS)N�error�{�G�z�?)
�	_continueryrz�popr{�status�
set_exception�
set_result�asyncio�sleep)r%�future�result�	cf_future�callback�typ�exc�tbs        r�_collectzDaskDistributedBackend._collect�s�����n�
	&�(,�(<�
%�
%�
%�
%�
%�
%�
%�n�f�f� �M�-�-�f�5�5�	��?�.�.�v�6�6���=�G�+�+�#)�L�C��b��+�+�C�0�0�0�0��(�(��0�0�0��H�V�$�$�$�$�)=��-��%�%�%�%�%�%�%�%�%��n�
	&�
	&�
	&�
	&�
	&s�Bc��tdfS)Nr=)rbr$s r�
__reduce__z!DaskDistributedBackend.__reduce__�s
��&��+�+r'c�0�t|j���dfS)N)rp���)rbrpr$s r�get_nested_backendz)DaskDistributedBackend.get_nested_backend�s��%�T�[�9�9�9�2�=�=r'rc�:�||_|�|��Sr!)�parallel�effective_n_jobs)r%�n_jobsr��backend_argss    r�	configurez DaskDistributedBackend.configure�s�� ��
��$�$�V�,�,�,r'c��d|_|jj�|j��t��|_dS)NT)r�rpre�add_callbackr�r�call_data_futuresr$s r�
start_callz!DaskDistributedBackend.start_call�s8��������%�%�d�m�4�4�4�!3�!5�!5����r'c�n�d|_tjd��|j���dS)NFr�)r��timer�r�r8r$s r�	stop_callz DaskDistributedBackend.stop_call�s8�����	
�
�4������$�$�&�&�&�&�&r'c	��t|j��������}|dks|js|S	|j�t���|j���nS#t$rF}d�	|jtdd|jz����}t|��|�d}~wwxYwt|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={})rc�)�sumrp�ncores�valuesrw�submitr`r��
_TimeoutError�format�maxr)r%r�r�r~�	error_msgs     rr�z'DaskDistributedBackend.effective_n_jobs�s���t�{�1�1�3�3�:�:�<�<�=�=���q� � ��(E� �#�#�
	1��K���1�2�2�9�9��5�
:�
7�
7�
7�
7���		1�		1�		1�H�
�f�T�2���Q��!>�>�?�?�A�A�

��y�)�)�q�0�����		1�����4�;�%�%�'�'�.�.�0�0�1�1�1s�	8B�
C�AC
�
Cc
�����K�t���t�dd������fd�}g}|jD]�\}}}t||���d{V����}tt	|���||������d{V������}|�|||f����t|��|fS)Nr�c���K�g}|D�]}t|��}|�vr|��|���2�	j�|d��}|�����	�|�d{V��}n#t$rYnwxYw|�`t|��rQt
|��dkr>�	j�|dd���}tj
|��}|�|<|�d{V��}|�|�|����|�|����|S)Ng@�@TF)�asynchronous�hash)r)rVrv�getr*rrrprtr��Task)
rH�out�arg�arg_idri�_coro�tr��itemgettersr%s
       ���r�maybe_to_futuresz>DaskDistributedBackend._to_func_args.<locals>.maybe_to_futuressV������C��(
$�(
$���C�����[�(�(��J�J�{�6�2�3�3�3���%�)�)�&�$�7�7���9�!2�!>��"3�C�"8�8�8�8�8�8�8����#�����������y�)�#�.�.�(�6�#�;�;��3D�3D�%)�K�$7�$7� #�-1�%*�%8�%�%�E�!(��U� 3� 3�A�56�-�c�2�&'�������A��=��J�J�q�M�M�M�M��J�J�s�O�O�O�O��Js�A)�)
A6�5A6)	�dict�getattr�itemsr@ru�keysr�rVrP)	r%rGr�rKrirHrIr�r�s	`      @@r�
_to_func_argsz$DaskDistributedBackend._to_func_argss��������f�f��$�D�*=�t�D�D��+	�+	�+	�+	�+	�+	�+	�Z��#�z�	,�	,�O�A�t�V��.�.�t�4�4�4�4�4�4�4�4�5�5�D��#�f�k�k�m�m�$4�$4�V�]�]�_�_�$E�$E�E�E�E�E�E�E�G�G�H�H�F��L�L�!�T�6�*�+�+�+�+��e���e�$�$r'c����tj�����j�_��fd�}�jj�|||���S)Nc��$�K���|���d{V��\}}t|���dt��j��}�jj|f||d��j��}�j�|��|�j	|<��j
|<dS)N�-)rKr0)r��reprr�hexrpr�rxry�addr{rz)rGr��batchrKr0�dask_futurer�r%s      ��rriz-DaskDistributedBackend.apply_async.<locals>.fGs������!%�!3�!3�D�!9�!9�9�9�9�9�9�9�L�E�5��%�[�[�0�0�5�7�7�;�0�0�C�,�$�+�,���"����/3�/A���K�
� �$�$�[�1�1�1�+3�D�O�K�(�)2�D�M�+�&�&�&r')�
concurrent�futures�Futurer�r�rprer�)r%rGr�rir�s`   @r�apply_asyncz"DaskDistributedBackend.apply_asyncBse�����&�-�-�/�/�	�!�(�	�
�		3�		3�		3�		3�		3�		3�	
���%�%�a��x�8�8�8��r'c�@�|jj5|jj���|jj���s<|jj���|jj����<ddd��dS#1swxYwYdS)z� Tell the client to cancel any task submitted via this instance

        joblib.Parallel will never access those results
        N)ry�lockr�r8�queue�emptyr�)r%�ensure_readys  r�abort_everythingz'DaskDistributedBackend.abort_everythingVs���
�
!�
&�	1�	1�� �(�.�.�0�0�0��*�0�6�6�8�8�
1��$�*�.�.�0�0�0��*�0�6�6�8�8�
1�	1�	1�	1�	1�	1�	1�	1�	1�	1�	1�	1�	1����	1�	1�	1�	1�	1�	1s�A9B�B�Bc#�K�ttd��rt��dV�ttd��rt��dSdS)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)�hasattrrrrr$s r�retrieval_contextz(DaskDistributedBackend.retrieval_context`sW�����<�!2�3�3�	��H�H�H�
�����<�!2�3�3�	��H�H�H�H�H�	�	r')NNNNrc)rNr!)T)r9r:r;�MIN_IDEAL_BATCH_DURATION�MAX_IDEAL_BATCH_DURATION�supports_timeoutr&r�r�r�r�r�r�r�r�r�r��
contextlib�contextmanagerr��
__classcell__)r�s@rrbrb�s�������"��"����48�BD�2�2�2�2�2�2�h&�&�&�,�,�,�>�>�>�-�-�-�-�6�6�6�
'�'�'�2�2�2�.;%�;%�;%�z����(1�1�1�1�����������r'rb),�
__future__rrrr��concurrent.futuresr�r�r��uuidrrr�rr	r
rrUrn�ImportError�
dask.utilsrr
�dask.sizeofr�dask.distributedrrrrrr�distributed.utilsrrr��tornado.genrrrCrNrPr`rbr=r'r�<module>r�s���@�@�@�@�@�@�@�@�@�@���������������������������J�J�J�J�J�J�J�J�J�J�&�&�&�&�&�&���K�K�K������������D��K�K�K��������/�/�/�/�/�/�/�/�/�"�"�"�"�"�"�����������������/�.�.�.�.�.�>�	D�C�C�C�C�C�C���>�>�>�=�=�=�=�=�=�=�=�>�������)�)�)�)�)�)�)�)�X���/�/�/���������.	�	�	�
a�a�a�a�a�.�0C�a�a�a�a�as!�?�	A�
A�7A>�>B�B

Youez - 2016 - github.com/yon3zu
LinuXploit