【发布时间】:2021-02-20 14:01:46
【问题描述】:
我正在开发一个爬虫(python 3.9),我需要启动 make_attempt() 并根据其执行结果启动应该同时工作的其他任务。
首先我创建初始任务并将其添加到存储所有异步任务的列表中:
self.data["worker"]["tasks"]
然后我启动:
await asyncio.gather(*self.data["worker"]["tasks"])
在 make_attempt() 中,我等待向服务器发出 POST 请求的结果(我使用 aiohttp 客户端),然后根据结果,我要么添加新任务,要么在一小段延迟后重复 make_attempt()。
我停止当前任务并将其从异步任务列表中删除,然后添加新任务。
async def make_attempt(self):
attempt: int = self.data["res"]["attempt"]
await self.do_something()
await sleep(1)
for task in self.data["worker"]["tasks"]:
print("Task name: %s" % task.get_name())
if task.get_name() == str(attempt):
task.cancel()
self.data["worker"]["tasks"] = [task for task in self.data["worker"]["tasks"] if task.get_name() != str(attempt)]
if 1 > 0: # a condition to start make_attempt() again
self.data["worker"]["tasks"].append(asyncio.create_task(self.make_attempt(), name=attempt))
await asyncio.gather(*self.data["worker"]["tasks"])
async def run(self):
self.data["worker"]["tasks"].append(asyncio.create_task(self.make_attempt(), name=self.data["res"]["attempt"]))
await asyncio.gather(*self.data["worker"]["tasks"])
我是 asyncio 的新手,所以也许您可以指出错误或建议更好的实现。
更新。这是我想要实现的架构:
main_task 应该运行,如果结果正常,它应该启动另一个任务的几个实例(参见 loop2)。当得到 loop2 中任务的结果时,应该运行一个新的子任务。 main_task 应该等待 loop2 中的所有任务完成或触发 Timeout。 loop2 中的所有任务都应该同时工作。
UPD 2. 此代码在 check_base_url() 方法执行大约 1000 个周期后生成 RecursionError: maximum recursion depth exceeded while calling a Python object。
class Scraper:
def __init__(self, data):
self.data = data
async def get_response(self, session, url, method="get", *args, **kwargs) -> Union[
ClientResponse, None]:
for _ in range(0, 20):
if method == "get":
try:
response = await session.get(url, headers={}, proxy="_proxy", *args, **kwargs)
if response.status > 399:
raise ScraperError(response.status)
await sleep(0.1)
return response
except (ClientError, ScraperError) as err:
await sleep(0.25)
continue
else:
try:
response = await session.post(url, headers={}, proxy="_proxy", *args, **kwargs)
if response.status > 399:
raise ScraperError(response.status)
await sleep(0.1)
return response
except (ClientError, ScraperError) as err:
await sleep(0.25)
continue
return None
async def get_captcha(self) -> SolvedCaptcha:
for _ in range(0, 20):
captcha = await self.task_1()
if captcha:
continue
async def final_task(self, url) -> bool:
async with ClientSession(cookies={}) as sess:
resp_step1: Union[ClientResponse, None] = await self.get_response(sess, "url", "post",
data={})
if resp_step1:
resp_step2: Union[ClientResponse, None] = await self.get_response(sess, "url", "get")
if resp_step2:
captcha: SolvedCaptcha = await self.get_captcha()
if captcha:
resp_captcha: Union[ClientResponse, None] = await self.get_response(sess, "url",
"post",
data={})
if resp_captcha:
if 2 > 1:
print("FINISHED")
return True
else:
return False
else:
return False
else:
return False
else:
return False
async def add_task_3(self) -> None:
if 2 > 1:
subtasks = [asyncio.create_task(self.final_task(self.data["res"]["slots_urls"][0]))]
await asyncio.gather(*subtasks)
else:
await self.add_task_3()
def parse(self, html: str, url: str) -> None:
soup = BeautifulSoup(html, "lxml")
# do parsing
async def task_2(self, url) -> bool:
async with ClientSession() as sess:
resp: Union[ClientResponse, None] = await self.get_response(sess, url)
if not resp:
return False
html = await resp.text()
self.parse(html, url)
async def add_task_2(self) -> None:
if 2 > 1:
subtasks = [asyncio.create_task(self.task_2(url)) for url in ["url1", "url2"]]
await asyncio.gather(*subtasks)
async def task_1(self) -> bool:
self.data["res"]["captcha_requested"] += 1
res = await self.captcha.task_1()
if not res:
return False
return True
async def add_task_1(self) -> None:
if 2 > 1:
subtasks = [asyncio.create_task(self.task_1()) for _ in range(0, 5)]
await asyncio.gather(*subtasks)
async def get_calendar_url(self, sess) -> bool:
resp: Union[ClientResponse, None] = await self.get_response(sess, "url", method="post",
data={})
if not resp:
return False
else:
return True
async def check_base_url(self) -> bool:
async with ClientSession() as session_0:
return await self.get_calendar_url(session_0)
async def schedule_tasks(self):
def start_again() -> bool:
if 2 > 1:
return True
return False
res_base_url: bool = await self.check_base_url()
if res_base_url:
tasks = [asyncio.create_task(self.add_task_1()),
asyncio.create_task(self.add_task_2()),
asyncio.create_task(self.add_task_3())]
await asyncio.gather(*tasks)
if start_again():
await sleep(0.1)
await self.schedule_tasks()
else:
await self.schedule_tasks()
async def run(self):
await self.schedule_tasks()
【问题讨论】:
-
有效吗?或者是什么问题?
-
您可以使用以
attempt为关键字的字典,而不是多次迭代您的任务列表。
标签: python python-asyncio python-3.8 python-3.9