o
    ôT·jÛ  ã                   @   s¬   d Z ddlmZ ddlmZ ddlmZ ddlZddl	Z	ddl
Zddl
mZ ddlmZ ddlmZ e	j d¡Ze	jje	jjgZe	j d	¡d
d„ ƒZdd„ Zedd„ ƒZdS )zd
Tests multithreading behaviour for reading and
parsing files for each parser defined in parsers.py
é    )Ú	ExitStack)ÚBytesIO)Ú
ThreadPoolN)Ú	DataFrame)ÚVersionÚpyarrow_xfailz0ignore:Passing a BlockManager:DeprecationWarningc                    sÔ   | }|j dkrt d¡}t|jƒtdƒk r| tjjdd�¡ d‰ d}‡ fdd„t|ƒD ƒ}t	ƒ �/‰‡fd	d
„|D ƒ}ˆ 
tdƒ¡}| |j|¡}|d }	|D ]}
t |	|
¡ qOW d   ƒ d S 1 scw   Y  d S )NÚpyarrowz16.0z+# ValueError: Found non-unique column index)Úreasonéd   é
   c                 3   s,   � | ]}d   dd„ tˆ ƒD ƒ¡ ¡ V  qdS )Ú
c                 S   s&   g | ]}|d ›d|d ›d|d ›�‘qS )Údú,© ©Ú.0Úir   r   úk/home/dinkstrade/pdmp-scanner/venv/lib/python3.10/site-packages/pandas/tests/io/parser/test_multi_thread.pyÚ
<listcomp>)   s   & zBtest_multi_thread_string_io_read_csv.<locals>.<genexpr>.<listcomp>N)ÚjoinÚrangeÚencode)r   Ú_)Úmax_row_ranger   r   Ú	<genexpr>(   s
   € ÿ
ÿz7test_multi_thread_string_io_read_csv.<locals>.<genexpr>c                    s   g | ]	}ˆ   t|ƒ¡‘qS r   )Úenter_contextr   )r   Úb)Ústackr   r   r   /   s    z8test_multi_thread_string_io_read_csv.<locals>.<listcomp>é   r   )ÚengineÚpytestÚimportorskipr   Ú__version__ÚapplymarkerÚmarkÚxfailr   r   r   r   ÚmapÚread_csvÚtmÚassert_frame_equal)Úall_parsersÚrequestÚparserÚpaÚ	num_filesÚbytes_to_dfÚfilesÚpoolÚresultsÚfirst_resultÚresultr   )r   r   r   Ú$test_multi_thread_string_io_read_csv   s*   

ÿ
þÿ"ør5   c                    sŒ   ‡‡fdd„}‡ ‡fdd„t ˆƒD ƒ}tˆd��}| ||¡}W d  ƒ n1 s)w   Y  |d j}|dd… D ]}	||	_q9t |¡}
|
S )	aš  
    Generate a DataFrame via multi-thread.

    Parameters
    ----------
    parser : BaseParser
        The parser object to use for reading the data.
    path : str
        The location of the CSV file to read.
    num_rows : int
        The number of rows to read per task.
    num_tasks : int
        The number of tasks to use for reading this DataFrame.

    Returns
    -------
    df : DataFrame
    c                    sB   | \}}|sˆ j ˆdd|dgd�S ˆ j ˆddt|ƒd |dgd�S )aj  
        Create a reader for part of the CSV.

        Parameters
        ----------
        arg : tuple
            A tuple of the following:

            * start : int
                The starting row to start for parsing CSV
            * nrows : int
                The number of rows to read.

        Returns
        -------
        df : DataFrame
        r   Údate)Ú	index_colÚheaderÚnrowsÚparse_datesNé   é	   )r7   r8   Úskiprowsr9   r:   )r'   Úint)ÚargÚstartr9   )r,   Úpathr   r   ÚreaderN   s   ÿ
úz0_generate_multi_thread_dataframe.<locals>.readerc                    s    g | ]}ˆ | ˆ ˆ ˆ f‘qS r   r   r   )Únum_rowsÚ	num_tasksr   r   r   p   s    ÿz4_generate_multi_thread_dataframe.<locals>.<listcomp>)Ú	processesNr   r;   )r   r   r&   ÚcolumnsÚpdÚconcat)r,   rA   rC   rD   rB   Útasksr1   r2   r8   ÚrÚfinal_dataframer   )rC   rD   r,   rA   r   Ú _generate_multi_thread_dataframe:   s   "ÿÿ

rL   c                 C   sð   d}d}| }d}t tj d¡ |¡tj d¡ |¡tj d¡ |¡tj d¡ |¡tj d¡ |¡dg| dg| dg| tjd|d	d
�tj|dd�dœ
ƒ}t |¡�}| 	|¡ t
||||ƒ}t ||¡ W d   ƒ d S 1 sqw   Y  d S )Né   é0   z__thread_pool_reader__.csvé   ÚfooÚbarÚbazz20000101 09:00:00Ús)ÚperiodsÚfreqÚint64)Údtype)
Úar   Úcr   ÚerP   rQ   rR   r6   r>   )r   ÚnpÚrandomÚdefault_rngrG   Ú
date_rangeÚaranger(   Úensure_cleanÚto_csvrL   r)   )r*   rD   rC   r,   Ú	file_nameÚdfrA   rK   r   r   r   Ú)test_multi_thread_path_multipart_read_csv€   s0   öÿ
ÿ"úrd   )Ú__doc__Ú
contextlibr   Úior   Úmultiprocessing.poolr   Únumpyr[   r    ÚpandasrG   r   Úpandas._testingÚ_testingr(   Úpandas.util.versionr   r$   ÚusefixturesÚxfail_pyarrowÚ
single_cpuÚslowÚ
pytestmarkÚfilterwarningsr5   rL   rd   r   r   r   r   Ú<module>   s&    þ

F