feat: 添加 GXGP Spider(广西政府采购网爬虫)+ 测试
This commit is contained in:
@@ -0,0 +1,117 @@
|
||||
import json
|
||||
import random
|
||||
import time
|
||||
from datetime import datetime
|
||||
from typing import List, Optional
|
||||
import httpx
|
||||
from app.config import settings
|
||||
from app.crawler.base import BaseSpider, CrawlResult, PipelineConfig
|
||||
from app.crawler.parsers import parse_gxgp_api_response, extract_pagination
|
||||
|
||||
|
||||
class GXGPSpider(BaseSpider):
|
||||
name = "gxgp"
|
||||
source_code = "gxgp"
|
||||
source_name = "广西政府采购网"
|
||||
|
||||
def __init__(self):
|
||||
self.base_url = settings.crawler_base_url
|
||||
self.announcement_api = f"{self.base_url}/portal/category"
|
||||
|
||||
def get_pipeline_config(self) -> PipelineConfig:
|
||||
return PipelineConfig(
|
||||
filter_enabled=True,
|
||||
keywords=list(settings.crawler_keywords),
|
||||
dedup_enabled=True,
|
||||
notify_mode="filtered",
|
||||
mark_sent=False,
|
||||
)
|
||||
|
||||
async def crawl(self, sources: Optional[List[str]] = None,
|
||||
max_pages: Optional[int] = None) -> CrawlResult:
|
||||
if max_pages is None:
|
||||
max_pages = settings.crawler_max_pages
|
||||
if sources is None:
|
||||
source_map = json.loads(settings.announcement_sources)
|
||||
sources = list(source_map.keys())
|
||||
|
||||
start_time = datetime.now()
|
||||
all_announcements = []
|
||||
error_messages = []
|
||||
|
||||
async with httpx.AsyncClient(timeout=settings.crawler_timeout) as client:
|
||||
for source_code in sources:
|
||||
source_info = json.loads(settings.announcement_sources).get(source_code)
|
||||
if not source_info:
|
||||
continue
|
||||
|
||||
category_id = source_info["category_id"]
|
||||
source_name = source_info["name"]
|
||||
|
||||
for page_no in range(1, max_pages + 1):
|
||||
if page_no > 1:
|
||||
await self._delay()
|
||||
|
||||
try:
|
||||
data = await self._fetch_page(
|
||||
client, source_code, category_id, page_no
|
||||
)
|
||||
if data is None:
|
||||
break
|
||||
|
||||
records = parse_gxgp_api_response(
|
||||
data, source_code, source_name,
|
||||
start_time, category_id
|
||||
)
|
||||
if not records:
|
||||
break
|
||||
|
||||
all_announcements.extend(records)
|
||||
|
||||
pagination = extract_pagination(data)
|
||||
if not pagination["has_next"]:
|
||||
break
|
||||
except Exception as e:
|
||||
error_messages.append(f"{source_code} page {page_no}: {e}")
|
||||
break
|
||||
|
||||
duration = (datetime.now() - start_time).total_seconds()
|
||||
return CrawlResult(
|
||||
source_code=self.source_code,
|
||||
source_name=self.source_name,
|
||||
total_count=len(all_announcements),
|
||||
new_count=len(all_announcements),
|
||||
announcements=all_announcements,
|
||||
error_message="; ".join(error_messages) if error_messages else None,
|
||||
crawled_at=start_time,
|
||||
duration=duration,
|
||||
)
|
||||
|
||||
async def _fetch_page(self, client: httpx.AsyncClient, source_code: str,
|
||||
category_id: int, page_no: int) -> Optional[dict]:
|
||||
payload = {
|
||||
"keyword": "",
|
||||
"publishDateBegin": "",
|
||||
"publishDateEnd": "",
|
||||
"pageNo": page_no,
|
||||
"pageSize": settings.crawler_page_size,
|
||||
"categoryCode": source_code,
|
||||
"_t": int(time.time() * 1000),
|
||||
}
|
||||
headers = {
|
||||
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
|
||||
"Content-Type": "application/json;charset=UTF-8",
|
||||
"Origin": self.base_url,
|
||||
"Referer": f"{self.base_url}/site/category?parentId={category_id}&childrenCode={source_code}",
|
||||
}
|
||||
response = await client.post(
|
||||
self.announcement_api, json=payload, headers=headers
|
||||
)
|
||||
if response.status_code != 200:
|
||||
return None
|
||||
return response.json()
|
||||
|
||||
async def _delay(self):
|
||||
import asyncio
|
||||
delay = random.uniform(1.0, 3.0)
|
||||
await asyncio.sleep(delay)
|
||||
Reference in New Issue
Block a user