ó »z�]c @sÞdZddlZddlZddlZddlZddlZddlZddlmZddl Z ddl m Z ddl mZddl mZddl mZddlmZdd lmZdd lmZdd lmZdd lmZdd lmZddlmZddlmZddlmZddlmZddlmZej e!ƒZ"dZ#ej$dddddddgƒZ%ej$ddddddddgƒZ&ej'd„ƒZ(d„Z)d e*fd!„ƒYZ+d"e*fd#„ƒYZ,d$efd%„ƒYZ-d&efd'„ƒYZ.d(e*fd)„ƒYZ/d*e*fd+„ƒYZ0d,e*fd-„ƒYZ1d.efd/„ƒYZ2e2j3d*e0ƒd0ej4fd1„ƒYZ5d2e5fd3„ƒYZ6d4e5fd5„ƒYZ7dS(6sCSpeeds up S3 throughput by using processes Getting Started =============== The :class:`ProcessPoolDownloader` can be used to download a single file by calling :meth:`ProcessPoolDownloader.download_file`: .. code:: python from s3transfer.processpool import ProcessPoolDownloader with ProcessPoolDownloader() as downloader: downloader.download_file('mybucket', 'mykey', 'myfile') This snippet downloads the S3 object located in the bucket ``mybucket`` at the key ``mykey`` to the local file ``myfile``. Any errors encountered during the transfer are not propagated. To determine if a transfer succeeded or failed, use the `Futures`_ interface. The :class:`ProcessPoolDownloader` can be used to download multiple files as well: .. code:: python from s3transfer.processpool import ProcessPoolDownloader with ProcessPoolDownloader() as downloader: downloader.download_file('mybucket', 'mykey', 'myfile') downloader.download_file('mybucket', 'myotherkey', 'myotherfile') When running this snippet, the downloading of ``mykey`` and ``myotherkey`` happen in parallel. The first ``download_file`` call does not block the second ``download_file`` call. The snippet blocks when exiting the context manager and blocks until both downloads are complete. Alternatively, the ``ProcessPoolDownloader`` can be instantiated and explicitly be shutdown using :meth:`ProcessPoolDownloader.shutdown`: .. code:: python from s3transfer.processpool import ProcessPoolDownloader downloader = ProcessPoolDownloader() downloader.download_file('mybucket', 'mykey', 'myfile') downloader.download_file('mybucket', 'myotherkey', 'myotherfile') downloader.shutdown() For this code snippet, the call to ``shutdown`` blocks until both downloads are complete. Additional Parameters ===================== Additional parameters can be provided to the ``download_file`` method: * ``extra_args``: A dictionary containing any additional client arguments to include in the `GetObject `_ API request. For example: .. code:: python from s3transfer.processpool import ProcessPoolDownloader with ProcessPoolDownloader() as downloader: downloader.download_file( 'mybucket', 'mykey', 'myfile', extra_args={'VersionId': 'myversion'}) * ``expected_size``: By default, the downloader will make a HeadObject call to determine the size of the object. To opt-out of this additional API call, you can provide the size of the object in bytes: .. code:: python from s3transfer.processpool import ProcessPoolDownloader MB = 1024 * 1024 with ProcessPoolDownloader() as downloader: downloader.download_file( 'mybucket', 'mykey', 'myfile', expected_size=2 * MB) Futures ======= When ``download_file`` is called, it immediately returns a :class:`ProcessPoolTransferFuture`. The future can be used to poll the state of a particular transfer. To get the result of the download, call :meth:`ProcessPoolTransferFuture.result`. The method blocks until the transfer completes, whether it succeeds or fails. For example: .. code:: python from s3transfer.processpool import ProcessPoolDownloader with ProcessPoolDownloader() as downloader: future = downloader.download_file('mybucket', 'mykey', 'myfile') print(future.result()) If the download succeeds, the future returns ``None``: .. code:: python None If the download fails, the exception causing the failure is raised. For example, if ``mykey`` did not exist, the following error would be raised .. code:: python botocore.exceptions.ClientError: An error occurred (404) when calling the HeadObject operation: Not Found .. note:: :meth:`ProcessPoolTransferFuture.result` can only be called while the ``ProcessPoolDownloader`` is running (e.g. before calling ``shutdown`` or inside the context manager). Process Pool Configuration ========================== By default, the downloader has the following configuration options: * ``multipart_threshold``: The threshold size for performing ranged downloads in bytes. By default, ranged downloads happen for S3 objects that are greater than or equal to 8 MB in size. * ``multipart_chunksize``: The size of each ranged download in bytes. By default, the size of each ranged download is 8 MB. * ``max_request_processes``: The maximum number of processes used to download S3 objects. By default, the maximum is 10 processes. To change the default configuration, use the :class:`ProcessTransferConfig`: .. code:: python from s3transfer.processpool import ProcessPoolDownloader from s3transfer.processpool import ProcessTransferConfig config = ProcessTransferConfig( multipart_threshold=64 * 1024 * 1024, # 64 MB max_request_processes=50 ) downloader = ProcessPoolDownloader(config=config) Client Configuration ==================== The process pool downloader creates ``botocore`` clients on your behalf. In order to affect how the client is created, pass the keyword arguments that would have been used in the :meth:`botocore.Session.create_client` call: .. code:: python from s3transfer.processpool import ProcessPoolDownloader from s3transfer.processpool import ProcessTransferConfig downloader = ProcessPoolDownloader( client_kwargs={'region_name': 'us-west-2'}) This snippet ensures that all clients created by the ``ProcessPoolDownloader`` are using ``us-west-2`` as their region. iÿÿÿÿN(tdeepcopy(tConfig(tMB(tALLOWED_DOWNLOAD_ARGS(tPROCESS_USER_AGENT(tMAXINT(t BaseManager(tCancelledError(tRetriesExceededError(tBaseTransferFuture(tBaseTransferMeta(tS3_RETRYABLE_DOWNLOAD_ERRORS(tcalculate_num_parts(tcalculate_range_parameter(tOSUtils(tCallArgstSHUTDOWNtDownloadFileRequestt transfer_idtbuckettkeytfilenamet extra_argst expected_sizet GetObjectJobt temp_filenametoffsetccs%tƒ}dVtjtj|ƒdS(N(t"_add_ignore_handler_for_interruptstsignaltSIGINT(toriginal_handler((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt ignore_ctrl_cs cCstjtjtjƒS(N(RRtSIG_IGN(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRstProcessTransferConfigcBs"eZdededd„ZRS(ii cCs||_||_||_dS(suConfiguration for the ProcessPoolDownloader :param multipart_threshold: The threshold for which ranged downloads occur. :param multipart_chunksize: The chunk size of each ranged download. :param max_request_processes: The maximum number of processes that will be making S3 API transfer-related requests at a time. N(tmultipart_thresholdtmultipart_chunksizetmax_request_processes(tselfR"R#R$((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt__init__s  (t__name__t __module__RR&(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR!stProcessPoolDownloadercBs­eZddd„Zddd„Zd„Zd„Zd„Zd„Zd„Z d„Z d„Z d „Z d „Z d „Zd „Zd „Zd„Zd„Zd„ZRS(cCs¸|dkri}nt|ƒ|_||_|dkrHtƒ|_ntjdƒ|_tjdƒ|_t ƒ|_ t |_ t jƒ|_d|_d|_d|_g|_dS(s­Downloads S3 objects using process pools :type client_kwargs: dict :param client_kwargs: The keyword arguments to provide when instantiating S3 clients. The arguments must match the keyword arguments provided to the `botocore.session.Session.create_client()` method. :type config: ProcessTransferConfig :param config: Configuration for the downloader ièN(tNonet ClientFactoryt_client_factoryt_transfer_configR!tmultiprocessingtQueuet_download_request_queuet _worker_queueRt_osutiltFalset_startedt threadingtLockt _start_lockt_managert_transfer_monitort _submittert_workers(R%t client_kwargstconfig((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR&#s         c CsÅ|jƒ|dkri}n|j|ƒ|jjƒ}td|d|d|d|d|d|ƒ}tjd|ƒ|jj |ƒt d|d|d|d|d|ƒ}|j ||ƒ} | S( ssDownloads the object's contents to a file :type bucket: str :param bucket: The name of the bucket to download from :type key: str :param key: The name of the key to download from :type filename: str :param filename: The name of a file to download to. :type extra_args: dict :param extra_args: Extra arguments that may be passed to the client operation :type expected_size: int :param expected_size: The expected size in bytes of the download. If provided, the downloader will not call HeadObject to determine the object's size and use the provided value instead. The size is needed to determine whether to do a multipart download. :rtype: s3transfer.futures.TransferFuture :returns: Transfer future representing the download RRRRRRs%Submitting download file request: %s.N( t_start_if_neededR*t_validate_all_known_argsR9tnotify_new_transferRtloggertdebugR0tputRt_get_transfer_future( R%RRRRRRtdownload_file_requestt call_argstfuture((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt download_fileDs"        cCs|jƒdS(shShutdown the downloader It will wait till all downloads are complete before returning. N(t_shutdown_if_needed(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pytshutdownqscCs|S(N((R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt __enter__xscGs?t|tƒr1|jdk r1|jjƒq1n|jƒdS(N(t isinstancetKeyboardInterruptR9R*tnotify_cancel_all_in_progressRJ(R%texc_typet exc_valuetargs((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt__exit__{scCs*|j�|js |jƒnWdQXdS(N(R7R4t_start(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR>�s  cCs+|jƒ|jƒ|jƒt|_dS(N(t_start_transfer_monitor_managert_start_submittert_start_get_object_workerstTrueR4(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRS†s   cCsCx<|D]4}|tkrtd|djtƒfƒ‚qqWdS(Ns/Invalid extra_args key '%s', must be one of: %ss, (Rt ValueErrortjoin(R%tprovidedtkwarg((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR?Œs   cCs1td|d|ƒ}td|jd|ƒ}|S(NRFRtmonitortmeta(tProcessPoolTransferMetatProcessPoolTransferFutureR9(R%RRFR]RG((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRD”s cCs?tjdƒtƒ|_|jjtƒ|jjƒ|_dS(Ns$Starting the TransferMonitorManager.(RARBtTransferMonitorManagerR8tstartRtTransferMonitorR9(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRT›s  c Cs`tjdƒtd|jd|jd|jd|jd|jd|jƒ|_ |j j ƒdS(Ns Starting the GetObjectSubmitter.ttransfer_configtclient_factoryttransfer_monitortosutiltdownload_request_queuet worker_queue( RARBtGetObjectSubmitterR-R,R9R2R0R1R:Ra(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRU¥s      c Cs~tjd|jjƒxat|jjƒD]M}td|jd|jd|jd|j ƒ}|j ƒ|j j |ƒq)WdS(NsStarting %s GetObjectWorkers.tqueueRdReRf( RARBR-R$trangetGetObjectWorkerR1R,R9R2RaR;tappend(R%t_tworker((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRV±s       cCs*|j�|jr |jƒnWdQXdS(N(R7R4t _shutdown(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRI¾s  cCs+|jƒ|jƒ|jƒt|_dS(N(t_shutdown_submittert_shutdown_get_object_workerst"_shutdown_transfer_monitor_managerR3R4(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRpÃs   cCstjdƒ|jjƒdS(Ns)Shutting down the TransferMonitorManager.(RARBR8RJ(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRsÉs cCs.tjdƒ|jjtƒ|jjƒdS(Ns%Shutting down the GetObjectSubmitter.(RARBR0RCtSHUTDOWN_SIGNALR:RY(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRqÍs cCsStjdƒx!|jD]}|jjtƒqWx|jD]}|jƒq;WdS(Ns#Shutting down the GetObjectWorkers.(RARBR;R1RCRtRY(R%RnRo((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRrÒs  N(R'R(R*R&RHRJRKRRR>RSR?RDRTRURVRIRpRsRqRr(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR)"s$! ,           R_cBs;eZd„Zed„ƒZd„Zd„Zd„ZRS(cCs||_||_dS(saThe future associated to a submitted process pool transfer request :type monitor: TransferMonitor :param monitor: The monitor associated to the proccess pool downloader :type meta: ProcessPoolTransferMeta :param meta: The metadata associated to the request. This object is visible to the requester. N(t_monitort_meta(R%R\R]((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR&Ûs cCs|jS(N(Rv(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR]èscCs|jj|jjƒS(N(Rutis_doneRvR(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pytdoneìscCsLy|jj|jjƒSWn+tk rG|jjƒ|jƒ‚nXdS(N(Rutpoll_for_resultRvRRMt_connecttcancel(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pytresultïs    cCs |jj|jjtƒƒdS(N(Rutnotify_exceptionRvRR(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR{s (R'R(R&tpropertyR]RxR|R{(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR_Ús    R^cBsDeZdZd„Zed„ƒZed„ƒZed„ƒZRS(s2Holds metadata about the ProcessPoolTransferFuturecCs||_||_i|_dS(N(t _transfer_idt _call_argst _user_context(R%RRF((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR& s  cCs|jS(N(R€(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRFscCs|jS(N(R(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRscCs|jS(N(R�(R%((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt user_contexts(R'R(t__doc__R&R~RFRR‚(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR^ s  R+cBseZdd„Zd„ZRS(cCs{||_|jdkr$i|_nt|jjdtƒƒƒ}|jsWt|_n|jdt7_||jdtt|ƒj|ƒ||_||_||_||_dS(süFulfills GetObjectJobs Downloads the S3 object, writes it to the specified file, and renames the file to its final location if it completes the final job for a particular transfer. :param queue: Queue for retrieving GetObjectJob's :param client_factory: ClientFactory for creating S3 clients :param transfer_monitor: Monitor for notifying :param osutil: OSUtils object to use for os-related behavior when performing the transfer. N(R©RlR&t_queueR,R9R2(R%RjRdReRf((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR&bs    cCsÂx»tr½|jjƒ}|tkr5tjdƒdS|jj|jƒsZ|j |ƒntjd|ƒ|jj |jƒ}tjd||jƒ|s|j |j|j |j ƒqqWdS(Ns Worker shutdown signal received.sBSkipping get object job %s because there was a previous exception.s%%s jobs remaining for transfer_id %s.(RWRÇR†RtRARBR9R˜Rt_run_get_object_jobR�t_finalize_downloadRR(R%tjobt remaining((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyR«us&     c Cs„y;|jd|jd|jd|jd|jd|jƒWnBtk r}tjd||dt ƒ|j j |j |ƒnXdS(NRRRRRsBException caught when downloading object for get object job %s: %sR®( t_do_get_objectRRRRRR°RARBRWR9R}R(R%RÊR±((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRÈŒs  c Cs¬d}x“t|jƒD]‚}y=|jjd|d||�}|j|||dƒdSWqtk r—} tjd| |d|jdt ƒ| }qXqWt |ƒ‚dS(NR·R¸tBodysCRetrying exception caught (%s), retrying request, (attempt %s / %s)iR®( R*Rkt _MAX_ATTEMPTSRªt get_objectt_write_to_fileR RARBRWR( R%RRRRRtlast_exceptionRÃtresponseR±((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRÌ™s   csbt|dƒ�M}|j|ƒt‡‡fd†dƒ}x|D]}|j|ƒqAWWdQXdS(Nsrb+csˆjˆjƒS(N(treadt _IO_CHUNKSIZE((tbodyR%(s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyt«tR×(topentseektitertwrite(R%RRRÕtftchunkstchunk((RÕR%s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRШs   cCsL|jj|ƒr%|jj|ƒn|j|||ƒ|jj|ƒdS(N(R9R˜R2t remove_filet_do_file_renameR“(R%RRR((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRɯscCsTy|jj||ƒWn6tk rO}|jj||ƒ|jj|ƒnXdS(N(R2t rename_fileR°R9R}Rß(R%RRRR±((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRà¶s ( R'R(RÎRRÔR&R«RÈRÌRÐRÉRà(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pyRl\s      (8Rƒt collectionst contextlibtloggingR.R5RtcopyRtbotocore.sessionRŠtbotocore.configRts3transfer.constantsRRRts3transfer.compatRRts3transfer.exceptionsRRts3transfer.futuresR R ts3transfer.utilsR R R RRt getLoggerR'RARtt namedtupleRRtcontextmanagerRRtobjectR!R)R_R^R+RbR‘R`tregistertProcessR¨RiRl(((s:/tmp/pip-build-kBFYxq/s3transfer/s3transfer/processpool.pytÂsp          ¸0`/t