【问题标题】:python - handling HTTP requests asynchronouslypython - 异步处理HTTP请求
【发布时间】:2016-05-04 14:16:32
【问题描述】:

我需要为 django 查询集的每个条目生成一个 PDF 报告。条目数量在 30k 到 40k 之间。

PDF 是通过外部 API 生成的。由于当前是按需生成的,因此通过 HTTP 请求/响应同步处理。 这个任务会有所不同,因为我想我将使用 django 管理命令来循环查询集并执行 PDF 生成。

我应该采用哪种方法来完成这项任务?我想到了 2 种可能的解决方案,虽然这些技术我从未使用过:

1) Celery:将任务(具有不同负载的http请求)分配给工作人员,然后在完成后检索它。

2)request-futures:以非阻塞方式使用请求。

目标是同时使用 API(例如同时发送 10 或 100 个 http 请求,具体取决于 API 可以处理多少并发请求)。

这里有没有人处理过类似的任务并可以就如何进行此任务提供建议?

以下是第一次尝试,使用multiprocessing(注意:大部分代码是重复使用的,不是我自己编写的,因为我拥有这个项目的所有权。):

class Checker(object):

    def __init__(self, *args, **kwargs):
        # ... various setup

    # other methods
    # .....

    def run_single(self, uuid, verbose=False):
        """
        run a single PDF generation and local download
        """
        start = timer()
        headers = self.headers

        data, obj = self.get_review_data(uuid)
        if verbose: 
            print("** Report: {} **".format(obj))
        response = requests.post(
            url=self.endpoint_url,
            headers=headers,
            data=json.dumps(data)
        )
        if verbose:
            print('POST - Response: {} \n {} \n {} secs'.format(
                response.status_code,
                response.content,
                response.elapsed.total_seconds())
            )
        run_url = self.check_progress(post_response=response, verbose=True)
        if run_url:
            self.get_file(run_url, obj, verbose=True)
        print("*** Download {}in {} secs".format("(verbose) " if verbose else "", timer()-start))


    def run_all(self, uuids, verbose=True):
        start = timer()
        for obj_uuid in review_uuids:
            self.run_single(obj_uuid, verbose=verbose)
        print("\n\n### Downloaded {}{} reviews in {} secs".format(
            "(verbose) " if verbose else "",
            len(uuids),
            timer() - start)
        )

    def run_all_multi(self, uuids, workers=4, verbose=True):
        pool = Pool(processes=workers)
        pool.map(self.run_single, uuids)


    def check_progress(self, post_response, attempts_limit=10000, verbose=False):
        """
        check the progress of PDF generation querying periodically the API endpoint
        """
        if post_response.status_code != 200:
            if verbose: print("POST response status code != 200 - exit")
            return None
        url = 'https://apidomain.com/{path}'.format(
            domain=self.domain,
            path=post_response.json().get('links', {}).get('self', {}).get('href'),
            headers = self.headers
        )
        job_id = post_response.json().get('jobId', '')
        status = 'Running'
        attempt_counter = 0
        start = timer()
        if verbose: 
            print("GET - url: {}".format(url))
        while status == 'Running':
            attempt_counter += 1
            job_response = requests.get(
                url=url,
                headers=self.headers,
            )
            job_data = job_response.json()
            status = job_data['status']
            message = job_data['message']
            progress = job_data['progress']
            if status == 'Error':
                if verbose:
                    end = timer()
                    print(
                        '{sc} - job_id: {job_id} - error_id: [{error_id}]: {message}'.format(
                            sc=job_response.status_code, 
                            job_id=job_id,
                            error_id=job_data['errorId'], 
                            message=message
                        ), '{} secs'.format(end - start)
                    )
                    print('Attempts: {} \n {}% progress'.format(attempt_counter, progress))
                return None
            if status == 'Complete':
                if verbose:
                    end = timer()
                    print('run_id: {run_id} - Complete - {secs} secs'.format(
                        run_id=run_id,
                        secs=end - start)
                    )
                    print('Attempts: {}'.format(attempt_counter))
                    print('{url}/files/'.format(url=url))
                return '{url}/files/'.format(url=url)
            if attempt_counter >= attempts_limit:
                if verbose:
                    end = timer()
                    print('File failed to generate after {att_limit} retrieve attempts: ({progress}% progress)' \
                          ' - {message}'.format(
                              att_limit=attempts_limit,
                              progress=int(progress * 100),
                              message=message
                          ), '{} secs'.format(end-start))
                return None
            if verbose:
                print('{}% progress  - attempts: {}'.format(progress, attempt_counter), end='\r')
                sys.stdout.flush()
            time.sleep(1)
        if verbose:
            end = timer()
            print(status, 'message: {} - attempts: {} - {} secs'.format(message, attempt_counter, end - start))
        return None

    def get_review_data(self, uuid, host=None, protocol=None):
        review = get_object_or_404(MyModel, uuid)
        internal_api_headers = {
            'Authorization': 'Token {}'.format(
                review.employee.csod_profile.csod_user_token
            )
        }

        data = requests.get(
            url=a_local_url,
            params={'format': 'json', 'indirect': 'true'},
            headers=internal_api_headers,
        ).json()
        return (data, review)

    def get_file(self, runs_url, obj, verbose=False):

        runs_files_response = requests.get(
            url=runs_url,
            headers=self.headers,
            stream=True,
        )

        runs_files_data = runs_files_response.json()


        file_path = runs_files_data['files'][0]['links']['file']['href'] # remote generated file URI
        file_response_url = 'https://apidomain.com/{path}'.format(path=file_path)
        file_response = requests.get(
            url=file_response_url,
            headers=self.headers,
            params={'userId': settings.CREDENTIALS['userId']},
            stream=True,
        )
        if file_response.status_code != 200:
            if verbose:
                print('error in retrieving file for {r_id}\nurl: {url}\n'.format(
                    r_id=obj.uuid, url=file_response_url)
                )
        local_file_path = '{temp_dir}/{uuid}-{filename}-{employee}.pdf'.format(
            temp_dir=self.local_temp_dir,
            uuid=obj.uuid,
            employee=slugify(obj.employee.get_full_name()),
            filename=slugify(obj.task.name)
        )
        with open(local_file_path, 'wb') as f:
            for block in file_response.iter_content(1024):
                f.write(block)
            if verbose:
                print('\n --> {r} [{uuid}]'.format(r=review, uuid=obj.uuid))
                print('\n --> File downloaded: {path}'.format(path=local_file_path))

    @classmethod
    def get_temp_directory(self):
        """
        generate a local unique temporary directory
        """
        return '{temp_dir}/'.format(
            temp_dir=mkdtemp(dir=TEMP_DIR_PREFIX),
        )

if __name__ == "__main__":
    uuids = #list or generator of objs uuids
    checker = Checker()
    checker.run_all_multi(uuids=uuids)

不幸的是,运行checker.run_all_multi有以下效果

  • python shell 冻结;
  • 不打印输出;
  • 没有生成文件;
  • 我必须从命令行终止控制台,正常的键盘中断停止工作

在运行checker.run_all 时,正常工作(一个接一个)。

关于为什么这段代码不起作用(而不是关于我可以使用什么来代替多处理)的任何建议?

谢谢大家。

【问题讨论】:

  • 您需要多久生成一次这些报告?生成是手动触发还是自动触发?
  • - 每年一次 - 手动
  • 在那个频率下,我倾向于使用 requests-futures 并避免需要设置 rabbitmq 等
  • @Anentropic 您能否更具体一些或提供任何示例/PoC?谢谢
  • Celery 可能更容易在 python 端编码,但您必须安装和配置其他软件(一个 MQ,例如 rabbit,加上运行 celery 工作进程)。 requests-futures 将是一个更简单的系统(只是一个 python 脚本),但您可能必须在 python 代码中进行某种速率限制,以避免向 PDF 服务发送 40k 请求

标签: python django httprequest python-requests python-multiprocessing


【解决方案1】:

根据您的频率,每年一次并手动进行。你不需要 Celery 或 request-futures。

创建一个类似的方法

def record_to_pdf(record):
    # create pdf from record

然后用代码创建一个管理命令(使用multiprocessing.Pool)

from multiprocessing import Pool
pool = Pool(processes=NUMBER_OF_CORES)
pool.map(record_to_pdf, YOUR_QUERYSET)

但管理命令不会是异步的。要使其异步,您可以在后台运行它。

此外,如果您的进程不受 CPU 限制(例如,它只是调用一些 API),那么正如 @Anentropic 建议的那样,您可以在创建池时尝试更多的进程。

【讨论】:

  • 对于不受 cpu 限制的任务,您还可以尝试处理的数量 > NUMBER_OF_CORES
  • @Anentropic 你是对的,record_to_pdf 方法只是调用一些 API,然后进程数可以增加很多(受网络速度和 API 速率限制)。
  • 试过了。它不会向标准输出输出任何内容,不会将任何文件保存到目标目录,也不会冻结 shell(我需要用 kill -9 杀死它)。相同的代码无需多处理即可工作,按顺序处理每个项目。我可以粘贴代码。有什么想法吗?
  • @Luke 请通过编辑您的问题发布代码。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2012-04-03
  • 2014-08-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多