Skip to content

Manage ClientConnectionResetError in Dataset._load_http_urls #279

Description

@github-actions

# TODO: Manage ClientConnectionResetError in Dataset._load_http_urls

        async def load_all(as_mime_type: None | str = None) -> 'Dataset[_ModelOrDatasetT]':
            tasks = []

            # TODO: Manage ClientConnectionResetError in Dataset._load_http_urls
            for host in hosts:
                host_config = self.config.http.for_host[host]
                async with RateLimitingClientSession(
                        host_config.requests_per_time_period,
                        host_config.time_period_in_secs) as client_session:
                    async with get_retry_client(
                            client_session=client_session,
                            retry_http_statuses=host_config.retry_http_statuses,
                            retry_attempts=host_config.retry_attempts,
                            retry_backoff_strategy=host_config.retry_backoff_strategy
                    ) as retry_client:
                        indices = hosts[host]
                        # fetch_task = get_auto_from_api_endpoint
                        # if as_mime_type:
                        #     match as_mime_type:
                        #         case 'application/json':
                        #             fetch_task = get_json_from_api_endpoint
                        #         case 'text/plain':
                        #             fetch_task = get_str_from_api_endpoint
                        #         case 'application/octet-stream' | _:
                        #             fetch_task = get_bytes_from_api_endpoint

                        ret = get_auto_from_api_endpoint.refine(
                            output_dataset_param='output_dataset').run(
                                http_url_dataset[indices],
                                retry_client=retry_client,
                                output_dataset=self,
                                as_mime_type=as_mime_type)

                        if not isinstance(ret, asyncio.Task):
                            assert inspect.iscoroutine(ret)
                            task = asyncio.create_task(ret)
                        else:
                            task = ret

                        tasks.append(task)

                        while not task.done():
                            await asyncio.sleep(ASYNC_LOAD_SLEEP_TIME)

            await asyncio.gather(*tasks)
            return self

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions