@summary: 分发任务 --------- @param tasks: --------- @result:
(self, tasks)
| 246 | return tasks |
| 247 | |
| 248 | def distribute_task(self, tasks): |
| 249 | """ |
| 250 | @summary: 分发任务 |
| 251 | --------- |
| 252 | @param tasks: |
| 253 | --------- |
| 254 | @result: |
| 255 | """ |
| 256 | if self._is_more_parsers: # 为多模版类爬虫,需要下发指定的parser |
| 257 | for task in tasks: |
| 258 | for parser in self._parsers: # 寻找task对应的parser |
| 259 | if parser.name in task: |
| 260 | if isinstance(task, dict): |
| 261 | task = PerfectDict(_dict=task) |
| 262 | else: |
| 263 | task = PerfectDict( |
| 264 | _dict=dict(zip(self._task_keys, task)), |
| 265 | _values=list(task), |
| 266 | ) |
| 267 | requests = parser.start_requests(task) |
| 268 | if requests and not isinstance(requests, Iterable): |
| 269 | raise Exception( |
| 270 | "%s.%s返回值必须可迭代" % (parser.name, "start_requests") |
| 271 | ) |
| 272 | |
| 273 | result_type = 1 |
| 274 | for request in requests or []: |
| 275 | if isinstance(request, Request): |
| 276 | request.parser_name = request.parser_name or parser.name |
| 277 | self._request_buffer.put_request(request) |
| 278 | result_type = 1 |
| 279 | |
| 280 | elif isinstance(request, Item): |
| 281 | self._item_buffer.put_item(request) |
| 282 | result_type = 2 |
| 283 | |
| 284 | if ( |
| 285 | self._item_buffer.get_items_count() |
| 286 | >= setting.ITEM_MAX_CACHED_COUNT |
| 287 | ): |
| 288 | self._item_buffer.flush() |
| 289 | |
| 290 | elif callable(request): # callbale的request可能是更新数据库操作的函数 |
| 291 | if result_type == 1: |
| 292 | self._request_buffer.put_request(request) |
| 293 | else: |
| 294 | self._item_buffer.put_item(request) |
| 295 | |
| 296 | if ( |
| 297 | self._item_buffer.get_items_count() |
| 298 | >= setting.ITEM_MAX_CACHED_COUNT |
| 299 | ): |
| 300 | self._item_buffer.flush() |
| 301 | |
| 302 | else: |
| 303 | raise TypeError( |
| 304 | "start_requests yield result type error, expect Request、Item、callback func, bug get type: {}".format( |
| 305 | type(requests) |
no test coverage detected