Files
surugaya/app/services/surugaya_client.py
T
2026-07-27 10:57:02 +08:00

975 lines
40 KiB
Python

"""骏河屋(suruga-ya.jp)抓取客户端
负责搜索列表和商品详情的抓取与解析。
通过 Cloudflare 会话复用机制,避免每次请求都触发验证挑战。
"""
from __future__ import annotations
import asyncio
import json
import logging
import re
from datetime import datetime
from pathlib import Path
from urllib.parse import parse_qs, urlencode, urlparse
from uuid import uuid4
from typing import Any
from selectolax.parser import HTMLParser
from app.core.config import Settings
from app.core.errors import ScrapeParseError, UpstreamBlockedError
from app.models.scrape import CategoryData, DetailRequest, OtherShopItem, OtherShopListData, ProductDetailData, \
ProductSummary, SearchRequest, SearchResultData, ShopAddress, ShopContactInfo, ShopDeliveryInfo, ShopInfoData, \
ShopInfoRequest, ShopItemsData, ShopItemsRequest, ShopPolicyInfo, ShopServiceInfo, ShopShippingFeeItem
from app.services.browser_pool import BrowserPool
from app.services.cloudflare_session import CloudflareSessionManager
from app.services.session_store import SessionStore
from app.utils.parse_util import format_price, get_shipping_fee, surugaya_photo_url_to_cdn
from surugaya_common.app_utils import is_empty_str
from surugaya_common.urls import BASE_URL, product_other_url, shop_section_url, shop_url
logger = logging.getLogger(__name__)
class SurugayaClient:
"""骏河屋抓取客户端,提供搜索和商品详情两个核心能力
工作流程:
1. 获取已通过 Cloudflare 验证的会话(优先复用,必要时新建)
2. 使用浏览器池中的页面,携带会话 Cookie 访问目标页面
3. 解析页面 HTML 提取结构化数据
"""
def __init__(
self,
settings: Settings,
browser_pool: BrowserPool,
session_manager: CloudflareSessionManager,
session_store: SessionStore,
):
self._settings = settings
self._browser_pool = browser_pool
self._session_manager = session_manager
self._session_store = session_store
self._categories_cache_key = "surugaya:categories:cache"
async def _raise_upstream_error(self, *, page: object, response: object | None, url: str) -> None:
status = getattr(response, "status", None)
try:
await self._session_manager.invalidate(url)
except Exception as exc:
logger.warning("清理上游异常会话失败:url=%s err=%s", url, exc)
screenshot_path: str | None = None
try:
screenshot_dir = Path(__file__).resolve().parents[2] / self._settings.log_dir / "screenshots"
screenshot_dir.mkdir(parents=True, exist_ok=True)
screenshot_file = screenshot_dir / f"{uuid4().hex}.png"
await page.screenshot(path=str(screenshot_file), full_page=True)
screenshot_path = str(screenshot_file)
except Exception:
screenshot_path = None
raise UpstreamBlockedError(f"Upstream status:{status}, 截图地址:{screenshot_path}")
async def search(self, payload: SearchRequest) -> SearchResultData:
"""抓取搜索列表页
流程:
1. 构建/使用搜索 URL
2. 获取已通过 Cloudflare 验证的会话
3. 用浏览器页面访问并获取 HTML
4. 解析商品总数和商品列表
Args:
payload: 搜索请求参数,包含 search_word、page 等
Returns:
SearchResultData: 包含搜索结果列表、总数、分页信息等
"""
logger.debug("开始抓取搜索页:payload=%s", json.dumps(payload.model_dump(), ensure_ascii=False))
url = str(payload.search_url) if payload.search_url else self._build_search_url(payload)
html, session = await self._session_manager.fetch_html(url)
tree = HTMLParser(html)
# 解析商品总数量,格式示例:該当件数:10,428件中 1-24件
current_page = self._extract_page_value(payload)
page_size = 24
total_count = self._extract_total_count(tree)
has_more = int(current_page * page_size < total_count) if total_count > 0 else 0
# 解析商品列表
items = self._parse_search_items(tree)
logger.info("搜索页解析完成 ===> 关键词: %s 抓取商品数=%s", payload.search_word, len(items))
return SearchResultData(
query=payload.search_word,
page=current_page,
page_size=page_size,
total_count=total_count,
has_more=has_more,
items=items,
session_id=session.session_id,
)
async def fetch_shop_items(self, payload: ShopItemsRequest) -> ShopItemsData:
"""抓取加盟店/市场店铺的商品列表
店铺商品列表复用站点的搜索页(/search?tenpo_code=xxx),
DOM 结构与普通搜索完全一致,因此解析逻辑与 search 共用。
Args:
payload: 店铺商品请求参数,包含 tenpo_cd、search_word、page 等
Returns:
ShopItemsData: 店铺名称、商品列表、总数和分页信息
"""
url = self._build_shop_items_url(payload)
logger.info("开始抓取店铺商品列表:tenpo_cd=%s page=%s", payload.tenpo_cd, payload.page)
logger.debug("店铺商品列表URL:url=%s", url)
html, _ = await self._session_manager.fetch_html(url)
tree = HTMLParser(html)
page_size = 24
total_count = self._extract_total_count(tree)
has_more = int(payload.page * page_size < total_count) if total_count > 0 else 0
# 列表页标题形如「駿河屋 山口大学前店の商品一覧」,去掉后缀即为店名
shop_name = ""
title_node = tree.css_first("#search_header .search_option .hit h2")
if title_node:
shop_name = re.sub(r"の商品一覧$", "", title_node.text().strip())
items = self._parse_search_items(tree)
logger.info("店铺商品列表解析完成 ===> tenpo_cd=%s 抓取商品数=%s", payload.tenpo_cd, len(items))
return ShopItemsData(
tenpo_cd=payload.tenpo_cd,
shop_name=shop_name,
query=payload.search_word,
page=payload.page,
page_size=page_size,
total_count=total_count,
has_more=has_more,
items=items,
)
async def fetch_shop_info(self, payload: ShopInfoRequest) -> ShopInfoData:
"""抓取加盟店/市场店铺信息
基础信息来自店铺主页 /shop/{tenpo_cd};当 include_details=True 时,
额外并发抓取 配送/ポリシー/返品保証/連絡 四个子页(多付出四次页面请求),
单个子页抓取或解析失败降级为 None,不影响主页信息返回。
Args:
payload: 店铺信息请求参数,包含 tenpo_cd、include_details
Returns:
ShopInfoData: 店铺信息
Raises:
ScrapeParseError: 店铺主页解析失败(如店铺不存在)
"""
tenpo_cd = payload.tenpo_cd
target_url = shop_url(tenpo_cd)
logger.info("开始抓取店铺信息:tenpo_cd=%s include_details=%s", tenpo_cd, payload.include_details)
html, _ = await self._session_manager.fetch_html(target_url)
info = self._parse_shop_info(HTMLParser(html), tenpo_cd)
if payload.include_details:
delivery, policy, service, contact = await asyncio.gather(
self._fetch_shop_section(tenpo_cd, "delivery", self._parse_shop_delivery),
self._fetch_shop_section(tenpo_cd, "policy", self._parse_shop_policy),
self._fetch_shop_section(tenpo_cd, "service", self._parse_shop_service),
self._fetch_shop_section(tenpo_cd, "contact", self._parse_shop_contact),
)
info.delivery = delivery
info.policy = policy
info.service = service
info.contact = contact
logger.debug("店铺信息解析完成:tenpo_cd=%s shop_name=%s", info.tenpo_cd, info.shop_name)
return info
async def _fetch_shop_section(self, tenpo_cd: str, section: str, parse: Any) -> Any:
"""抓取并解析单个店铺子页;抓取或解析失败时降级为 None,不影响主流程。"""
section_url = shop_section_url(section, tenpo_cd)
try:
html, _ = await self._session_manager.fetch_html(section_url)
except Exception as exc:
logger.warning("抓取店铺子页失败:tenpo_cd=%s section=%s err=%s", tenpo_cd, section, exc)
return None
try:
return parse(HTMLParser(html))
except Exception:
logger.exception("解析店铺子页失败:tenpo_cd=%s section=%s", tenpo_cd, section)
return None
@staticmethod
def _parse_shop_info(tree: HTMLParser, tenpo_cd: str) -> ShopInfoData:
"""解析店铺主页 /shop/{tenpo_cd} 的基础信息(店名、logo、评分、公告)。"""
brand = tree.css_first(".shop_brand")
if brand is None:
raise ScrapeParseError(f"店铺信息解析失败:tenpo_cd={tenpo_cd}")
name_node = brand.css_first(".shop_info h1")
shop_name = name_node.text().strip() if name_node else ""
if not shop_name:
raise ScrapeParseError(f"店铺名称解析失败:tenpo_cd={tenpo_cd}")
logo_node = brand.css_first(".shop_logo img")
logo_src = (logo_node.attributes.get("src") or "") if logo_node else ""
logo_url = SurugayaClient._normalize_url(logo_src) if logo_src else ""
# 评分区形如:<div class="pull-left padR20">5.0</div><div class="pull-right">(1345件)</div>
rating_score = ""
rating_count = 0
point_node = brand.css_first(".shop_info .point")
if point_node:
score_node = point_node.css_first(".pull-left")
if score_node:
score_match = re.search(r"[\d.]+", score_node.text())
if score_match:
rating_score = score_match.group(0)
count_node = point_node.css_first(".pull-right")
if count_node:
count_match = re.search(r"(\d+)", count_node.text().replace(",", ""))
if count_match:
rating_count = int(count_match.group(1))
notice_node = tree.css_first("#search_result .contact-top .content_text")
notice = SurugayaClient._clean_block_text(notice_node)
return ShopInfoData(
tenpo_cd=tenpo_cd,
shop_name=shop_name,
shop_url=shop_url(tenpo_cd),
logo_url=logo_url,
rating_score=rating_score,
rating_count=rating_count,
notice=notice,
)
@staticmethod
def _parse_shop_delivery(tree: HTMLParser) -> ShopDeliveryInfo:
"""解析店铺「配送」子页:发货说明 + 按都道府县的运费表。"""
description = SurugayaClient._clean_block_text(tree.css_first("#main2 .content_text"))
shipping_fee_list: list[ShopShippingFeeItem] = []
for row in tree.css("#main2 table.table-bordered tr"):
tds = row.css("td")
if len(tds) < 2:
# 表头行(th)无 td,跳过
continue
prefecture = tds[0].text().strip()
fee = format_price(tds[1].text().strip()) or 0
if prefecture:
shipping_fee_list.append(ShopShippingFeeItem(prefecture=prefecture, fee=fee))
return ShopDeliveryInfo(description=description, shipping_fee_list=shipping_fee_list)
@staticmethod
def _parse_shop_policy(tree: HTMLParser) -> ShopPolicyInfo:
"""解析店铺「ポリシー」子页:政策全文(含特定商取引法表记)。"""
return ShopPolicyInfo(description=SurugayaClient._clean_block_text(tree.css_first("#main2 .content_text")))
@staticmethod
def _parse_shop_service(tree: HTMLParser) -> ShopServiceInfo:
"""解析店铺「返品、保証、払い戻し」子页。
页面结构为 h4 标题与 content_text 正文交替排列,按出现顺序配对,
再依据标题文案分派到 return_policy / warranty。
"""
info = ShopServiceInfo()
container = tree.css_first("#main2 #search_result > div")
if container is None:
return info
current_title = ""
for node in container.iter():
if node.tag == "h4":
current_title = node.text().strip()
continue
if "content_text" not in (node.attributes.get("class") or ""):
continue
text = SurugayaClient._clean_block_text(node)
if "保証" in current_title:
info.warranty = text
elif "返品" in current_title or "払い戻し" in current_title:
info.return_policy = text
return info
@staticmethod
def _parse_shop_contact(tree: HTMLParser) -> ShopContactInfo:
"""解析店铺「連絡」子页:联系地址与返品地址。"""
field_map = {
"郵便番号": "postal_code",
"都道府県": "prefecture",
"市区町村": "city",
"番地": "street",
"ビル・マンション名": "building",
}
address_list: list[ShopAddress] = []
for block in tree.css("#main2 .contact-top"):
label_node = block.css_first("h4")
if label_node is None:
continue
address = ShopAddress(label=label_node.text().strip())
for p in block.css("p"):
key, sep, value = p.text().partition(":")
if not sep:
continue
field = field_map.get(key.strip())
if field:
setattr(address, field, value.strip())
if address.postal_code or address.prefecture:
address_list.append(address)
return ShopContactInfo(address_list=address_list)
@staticmethod
def _clean_block_text(node: object | None) -> str:
"""提取节点文本并归一化空白:去掉行首尾空白与连续空行,保留段落换行。"""
if node is None:
return ""
lines = [line.strip() for line in node.text().splitlines()]
return "\n".join(line for line in lines if line)
@staticmethod
def _extract_total_count(tree: HTMLParser) -> int:
"""从搜索/店铺列表页头部解析命中总数,格式示例:該当件数:10,428件中 1-24件"""
hit_node = tree.css_first("#search_header .search_option .hit")
if hit_node is None:
return 0
match = re.search(r"該当件数:([\d,]+)件", hit_node.text())
return int(match.group(1).replace(",", "")) if match else 0
async def fetch_detail(self, payload: DetailRequest) -> ProductDetailData:
"""抓取商品详情页
流程:
1. 根据请求参数确定商品 URL 和 ID(支持直接传 ID 或完整 URL)
2. 获取已通过 Cloudflare 验证的会话
3. 用浏览器页面访问并获取 HTML
4. 从 DOM 中提取商品名称、图片、价格、售卖状态、描述等
Args:
payload: 详情请求参数,包含 id 或 product_url
Returns:
ProductDetailData: 商品详情数据
Raises:
ScrapeParseError: 当 id 和 product_url 都未提供时
"""
goods_id = payload.id or ""
if is_empty_str(goods_id):
raise ScrapeParseError("id is empty")
url = self._build_product_url(goods_id, tenpo_cd=payload.tenpo_cd, branch_number=payload.branch_number)
logger.info("开始抓取详情页:goods_id=%s", goods_id)
logger.debug("详情页URL:url=%s", url)
# 抓取页面 HTML(会话复用失败时由 fetch_html 内部自动重建,避免成功/失败交替)
html, _ = await self._session_manager.fetch_html(url)
tree = HTMLParser(html)
# 提取商品标题
title_node = tree.css_first("#item_title")
goods_name = title_node.text().strip() if title_node else None
if not goods_name:
raise ScrapeParseError("商品名称解析失败")
# 提取商品主图
img_node = tree.css_first("div.easyzoom--with-thumbnails img")
goods_image: str = (img_node.attributes.get("src") or "") if img_node else ""
if goods_image and "database/images/no_photo.jpg" in goods_image:
goods_image = "https://oss.daimatech.com/suruga-ya/no_photo.jpg"
# 提取商品价格:优先从 .price_choise 中取,其次从 span.purchase-price 中取
price_text = ""
price_choise = tree.css_first("#cart .price_choise")
if price_choise:
buy_node = price_choise.css_first(".text-price-detail.price-buy")
if buy_node:
price_text = buy_node.text().strip()
if not price_text:
purchase_node = tree.css_first("span.purchase-price")
if purchase_node:
price_text = purchase_node.text().strip()
goods_price: int = format_price(price_text) or 0 if price_text else 0
# 判断售卖状态:
out_of_stock_dom = tree.css_first("div.out-of-stock-text")
on_sold = 0 if out_of_stock_dom else 1
# 解析规格列表
sku = ""
sku_list = []
pdom = tree.css_first("#item_title")
if pdom and pdom.parent:
direct_children = pdom.parent.css(".price_group > .item-price")
sku_index = 0
for node in direct_children:
sku_index += 1
vo: dict[str, Any] = {"sku_name": "", "price": 0, "stock": 0}
label_node = node.css_first("label")
if label_node:
vo["sku_name"] = label_node.attributes.get("data-label", "").strip()
input_node = node.css_first("input.form-check-input")
if input_node:
zaiko_str = input_node.attributes.get("data-zaiko", "")
if zaiko_str and zaiko_str != "null":
try:
zaiko = json.loads(zaiko_str)
if isinstance(zaiko, dict):
vo.update(zaiko)
vo["price"] = int(zaiko.get("baika", 0))
vo["stock"] = int(zaiko.get("zaiko", 0))
vo["sku_id"] = str(sku_index)
if sku_index == 1:
sku = vo["sku_id"]
except json.JSONDecodeError:
vo["zaiko_raw"] = zaiko_str
sku_list.append(vo)
# 提取库存
stock = 0
if sku_list:
# 从 SKU 列表中找到匹配当前 SKU 的库存
stock = next((item["stock"] for item in sku_list if item.get("sku_id") == sku), 0)
else:
quantity_selection = tree.css_first("#quantity_selection")
if quantity_selection:
stock = 1
last_option = quantity_selection.css_first("option:last-child")
if last_option:
option_val = last_option.attributes.get("value")
if option_val and option_val.isdigit():
stock = int(option_val)
# 提取商品描述
desc_note = tree.css_first("p.note.text-break")
description = desc_note.text().strip() if desc_note else ""
# 提取店铺来源
supplier = "" # 默认表示官方自营
# 利用与目标 div.mb-2 平级的 #naviplus-review-list-5 缩小查找范围以提升效率
review_node = tree.css_first("#naviplus-review-list-5")
if review_node and review_node.parent:
context_node = review_node.parent
for node in context_node.css("div.mb-2"):
text = node.text()
if "この商品は" in text and "販売" in text and "発送" in text:
# 提取 mb-2 下的全部文本,并清理多余的换行和空白字符
supplier = " ".join(text.split())
break
# 提取价格备注
price_note = None
price_note_node = tree.css_first(".item-price-note-wrap")
if price_note_node:
price_note = price_note_node.text().strip()
# 提取运费相关
shipping_comission_fee = get_shipping_fee(goods_price)
# 提取面包屑
breadcrumb_list = self.get_breadcrumb_list(tree)
# 仅当调用方显式要求时,才额外抓取「他のショップ」页面(多付出一次 Cloudflare 保护页面请求)
other_shop_list: list[OtherShopItem] = []
if payload.include_other_shops:
other_shop_list = await self._fetch_other_shop_list(goods_id)
detail = ProductDetailData(
goods_id=goods_id,
goods_name=goods_name,
goods_image=goods_image,
goods_price=goods_price,
stock=stock,
on_sold=on_sold,
sku=sku,
sku_list=sku_list,
description=description,
supplier=supplier,
price_note=price_note,
shipping_comission_fee=shipping_comission_fee,
breadcrumb_list=breadcrumb_list,
other_shop_list=other_shop_list,
)
logger.debug("详情页解析完成:goods_id=%s goods_name=%s ", detail.goods_id, detail.goods_name)
return detail
async def fetch_other_shops(self, payload: DetailRequest) -> OtherShopListData:
"""获取商品「他のショップ」全部店铺报价列表
独立于 fetch_detail,仅在需要展示/跳转其他店铺时按需调用,
避免每次详情请求都多付出一次 Cloudflare 保护页面抓取的成本。
Args:
payload: 详情请求参数,仅使用其中的 id
Returns:
OtherShopListData: 全部店铺(含官方自营+加盟店/市场店铺)报价列表
Raises:
ScrapeParseError: 当 id 未提供时
"""
goods_id = payload.id or ""
if is_empty_str(goods_id):
raise ScrapeParseError("id is empty")
other_shop_list = await self._fetch_other_shop_list(goods_id)
return OtherShopListData(goods_id=goods_id, other_shop_list=other_shop_list)
async def _fetch_other_shop_list(self, goods_id: str) -> list[OtherShopItem]:
"""抓取并解析商品「他のショップ」页面,返回全部店铺(含官方自营)报价列表
该页面独立于详情页,需额外一次 Cloudflare 保护页面请求;
解析失败或抓取异常时降级为空列表,不影响详情页主流程。
"""
other_url = product_other_url(goods_id)
try:
other_html, _ = await self._session_manager.fetch_html(other_url)
except Exception as exc:
logger.warning("抓取他のショップ页面失败:goods_id=%s err=%s", goods_id, exc)
return []
try:
return self._parse_other_shop_list(HTMLParser(other_html))
except Exception:
logger.exception("解析他のショップ页面失败:goods_id=%s", goods_id)
return []
@staticmethod
def _parse_other_shop_list(tree: HTMLParser) -> list[OtherShopItem]:
"""解析「他のショップ」页面的全部店铺(tbl_all 标签页)报价列表
每行对应一个店铺报价;tenpo_cd 为空表示官方自营(駿河屋本店)。
"""
items: list[OtherShopItem] = []
for row in tree.css("#tbl_all tr.item"):
tds = row.css("td")
if len(tds) < 3:
continue
price_node = tds[0].css_first("strong.text-red")
price = format_price(price_node.text().strip()) or 0 if price_node else 0
condition_node = tds[1].css_first("h2.title_product")
condition = condition_node.text().strip() if condition_node else ""
detail_link_node = tds[1].css_first("a[href*='/product/detail/']")
detail_href = detail_link_node.attributes.get("href") if detail_link_node else ""
detail_url = f"{BASE_URL}{detail_href}" if detail_href else ""
shop_col = tds[2]
shop_link_node = shop_col.css_first("a[href^='/shop/']")
tenpo_cd = ""
shop_name = "駿河屋"
shop_url = f"{BASE_URL}/"
if shop_link_node:
shop_href = shop_link_node.attributes.get("href") or ""
tenpo_cd = shop_href.rsplit("/", 1)[-1]
shop_name = shop_link_node.text().strip()
shop_url = f"{BASE_URL}{shop_href}"
rating_score = ""
rating_count = 0
rating_node = shop_col.css_first("div:not([class])")
if rating_node:
rating_text = rating_node.text().strip()
rating_match = re.match(r"([\d.]+)\((\d+)", rating_text)
if rating_match:
rating_score = rating_match.group(1)
rating_count = int(rating_match.group(2))
shipping_note_node = row.css_first("ul.shipping_campaing_info li.padT5")
shipping_note = shipping_note_node.text().strip() if shipping_note_node else ""
# branch_number 优先从 detail_url 的查询参数中提取(与 detail_url 保持一致),
# data-zaiko_data 仅作为兜底(极少数情况下 href 未携带 branch_number 时)。
branch_number = ""
if detail_href:
branch_number = (parse_qs(urlparse(detail_href).query).get("branch_number") or [""])[0]
if not branch_number:
campaign_node = row.css_first("li.ajax-campaign-placeholder")
if campaign_node:
zaiko_raw = campaign_node.attributes.get("data-zaiko_data") or ""
if zaiko_raw:
try:
zaiko = json.loads(zaiko_raw)
branch_number = str(zaiko.get("branch_number", ""))
except json.JSONDecodeError:
pass
items.append(
OtherShopItem(
tenpo_cd=tenpo_cd,
shop_name=shop_name,
shop_url=shop_url,
is_official=tenpo_cd == "",
condition=condition,
price=price,
branch_number=branch_number,
rating_score=rating_score,
rating_count=rating_count,
shipping_note=shipping_note,
detail_url=detail_url,
)
)
return items
async def fetch_categories(self) -> list[CategoryData]:
"""获取商品分类(带缓存)
解析骏河屋的商品分类,由于分类不易变动,
优先从 redis 获取缓存,若 redis 未配置或缓存失效,则重新抓取。
默认缓存 24 小时。
"""
# 如果配置了 Redis,尝试从 Redis 获取缓存
if self._session_store._redis is not None:
try:
cached_data = await self._session_store._redis.get(self._categories_cache_key)
if cached_data:
logger.debug("命中商品分类 Redis 缓存")
data_list = json.loads(cached_data)
return [CategoryData(**item) for item in data_list]
except Exception as e:
logger.warning("读取 Redis 缓存失败: %s", e)
logger.info("开始获取并解析商品分类...")
top_categories: list[CategoryData] = []
target_url = f"{BASE_URL}/"
session = await self._session_manager.get_verified_session(target_url, True)
html_content = ""
if session.html is not None:
html_content = session.html
else:
async with self._browser_pool.page_session(storage_state=session.storage_state) as (_, page):
try:
response = await page.goto(
target_url,
wait_until="domcontentloaded",
timeout=int(self._settings.request_timeout_seconds * 1000),
)
status = response.status if response is not None else None
if status != 200:
raise ScrapeParseError(f"抓取首页分类失败,状态码 {status}")
html_content = await page.content()
except Exception as e:
logger.error("抓取首页分类失败:%s", e)
raise ScrapeParseError(f"抓取首页分类失败: {e!s}") from e
tree = HTMLParser(html_content)
# 解析一级分类列表
cate_nodes = tree.css(".catemenu_group .cate_menu .cate_item")
for cate_item in cate_nodes:
href = cate_item.css_first("a").attributes.get("href")
if not href:
continue
# 提取 href 的路径名(不含 .html 和前置的 /,例如:/avsoft.html -> avsoft)
match = re.search(r'/([^/]+)(?:\.html|/)', href)
if not match:
continue
cate_id = match.group(1)
# 过滤掉指定的分类
if cate_id in ["boyslove"]:
continue
# 提取包含的文字内容(拼接多个 span,例如 おもちゃ・ + ホビー)
h4_node = cate_item.css_first("a h4")
if not h4_node:
continue
cate_name = h4_node.text().strip()
# 清理多余空格与换行,防止拼接文本里出现多余空白
cate_name = "".join(cate_name.split())
icon_node = cate_item.css_first("a img")
icon = icon_node.attributes.get("src") if icon_node else None
pid = "0"
top_categories.append(CategoryData(id=cate_id, name=cate_name, icon=icon, pid=pid, href=href))
# 解析二级分类列表
seen_pairs: set[tuple[str, str]] = set()
async with self._browser_pool.page_session(storage_state=session.storage_state) as (_, page):
for top_category in top_categories:
try:
category_url = f"{BASE_URL}{top_category.href}"
response = await page.goto(
category_url,
wait_until="domcontentloaded",
timeout=int(self._settings.request_timeout_seconds * 1000),
)
status = response.status if response is not None else None
if status != 200:
await self._raise_upstream_error(page=page, response=response, url=category_url)
category_html = await page.content()
category_tree = HTMLParser(category_html)
menu_node = category_tree.css_first("#menu")
if not menu_node:
continue
related_block = None
for block in menu_node.css(".block"):
h2_node = block.css_first("h2")
if h2_node and "関連ジャンルで絞り込む" in h2_node.text():
related_block = block
break
if not related_block:
continue
for a_node in related_block.css("ul.border_bottom li a"):
sub_name = "".join(a_node.text().split())
if not sub_name:
continue
sub_href = a_node.attributes.get("href") or ""
parsed = urlparse(sub_href)
sub_id = ""
if parsed.path == "/search":
query = parse_qs(parsed.query)
sub_id = (query.get("category") or [""])[0]
if not sub_id:
match = re.search(r"/([^/]+)\.html$", parsed.path)
if match:
sub_id = match.group(1)
if not sub_id:
continue
pair = (top_category.id, sub_id)
if pair in seen_pairs:
continue
seen_pairs.add(pair)
top_category.children.append(
CategoryData(id=sub_id, name=sub_name, pid=top_category.id, href=sub_href))
# 休眠1秒,避免对服务器造成过大压力
await asyncio.sleep(1)
except Exception as e:
logger.warning("抓取二级分类失败:top_id=%s err=%s", top_category.id, e)
continue
# 更新 Redis 缓存
if top_categories and self._session_store._redis is not None:
try:
logger.debug("更新商品分类 Redis 缓存")
data_list = [item.model_dump() for item in top_categories]
# 缓存有效期 24 小时 (86400秒)
await self._session_store._redis.setex(
self._categories_cache_key,
86400 * 30,
json.dumps(data_list, ensure_ascii=False)
)
except Exception as e:
logger.warning("写入 Redis 缓存失败: %s", e)
return top_categories
# ------------------------------------------------------------------
# 内部辅助方法
# ------------------------------------------------------------------
def _build_search_url(self, payload: SearchRequest) -> str:
"""根据搜索参数构建骏河屋搜索 URL
将 payload 中的 search_word、page 等参数拼接为完整的搜索地址。
search_url 字段本身不参与拼接,它用于直接指定完整 URL。
"""
params = payload.model_dump(exclude_none=True)
params.pop("search_url", None)
# if not str(params.get("search_word", "")).strip():
# raise ScrapeParseError("search_word is required")
query_string = urlencode(params, doseq=True)
base_url = self._settings.base_url.rstrip("/")
return f"{base_url}/search?{query_string}"
@staticmethod
def _extract_page_value(payload: SearchRequest) -> int:
"""安全地提取页码,确保返回 >= 1 的整数"""
page_raw = payload.model_dump(exclude_none=True).get("page", 1)
try:
current_page = int(page_raw)
except (TypeError, ValueError):
return 1
return current_page if current_page >= 1 else 1
def _build_shop_items_url(self, payload: ShopItemsRequest) -> str:
"""构建店铺商品列表 URL
店内商品列表即站点搜索页按 tenpo_code 过滤的结果;
tenpo_cd 转成站点侧的 tenpo_code 参数,其余字段(rankBy、category 等)原样透传。
"""
params = payload.model_dump(exclude_none=True)
params["tenpo_code"] = params.pop("tenpo_cd")
query_string = urlencode(params, doseq=True)
base_url = self._settings.base_url.rstrip("/")
return f"{base_url}/search?{query_string}"
def _build_product_url(
self,
product_id: str | None,
tenpo_cd: str | None = None,
branch_number: str | None = None,
) -> str:
"""根据商品 ID 构建详情页 URL;tenpo_cd/branch_number 用于跳转查看指定店铺的报价"""
if not product_id:
raise ScrapeParseError("Either product_id or product_url must be provided")
base_url = self._settings.base_url.rstrip("/")
url = f"{base_url}/product/detail/{product_id}"
params = []
if tenpo_cd:
params.append(f"tenpo_cd={tenpo_cd}")
if branch_number:
params.append(f"branch_number={branch_number}")
if params:
url += "?" + "&".join(params)
return url
def _parse_search_items(self, tree: HTMLParser) -> list[ProductSummary]:
"""解析搜索结果列表
从搜索页 DOM 中提取每个商品的 ID、名称、图片、价格、市场价、售卖状态。
自动去重(按 goods_id),最多返回 50 条结果。
"""
item_nodes = tree.css("div.item_box div.item")
if not item_nodes:
item_nodes = tree.css("div.item")
items: list[ProductSummary] = []
seen_goods_ids: set[str] = set()
for node in item_nodes:
link_node = node.css_first("div.title a") or node.css_first("div.photo_box a")
href = link_node.attributes.get("href") if link_node else None
if not href:
continue
goods_link = self._normalize_url(href)
goods_id = self._extract_product_id(goods_link) or self._extract_product_id(href)
if not goods_id:
continue
if goods_id in seen_goods_ids:
continue
seen_goods_ids.add(goods_id)
name_node = node.css_first("h3.product-name")
goods_name = name_node.text().strip() if name_node else f"product-{goods_id}"
img_node = node.css_first("div.photo_box img")
goods_image_str = (img_node.attributes.get("src") or "").strip() if img_node else ""
goods_image = surugaya_photo_url_to_cdn(goods_image_str, goods_id)
# 解析商品售价(自营价格)
goods_price = 0
price_teika = node.css_first(".item_price .price_teika")
if price_teika is not None:
strong = price_teika.css_first(".text-red strong")
price_text = strong.text().strip() if strong is not None else price_teika.text().strip()
parsed_goods_price = format_price(price_text)
if parsed_goods_price is not None:
goods_price = parsed_goods_price
# 解析第三方市场价
third_shop_price = 0
makepla_strong = node.css_first(".item_price .makeplaTit .text-red strong")
if makepla_strong is not None:
price_text = makepla_strong.text().strip()
parsed_third_shop_price = format_price(price_text)
if parsed_third_shop_price is not None:
third_shop_price = parsed_third_shop_price
# 判断售卖状态:品切れ = 售罄
on_sold = 1
sold_out = node.css_first(".item_price p.price")
if sold_out is not None and "品切れ" in sold_out.text().strip():
on_sold = 0
# 售罄时,售价回退到市场价
goods_price = third_shop_price
now = str(int(datetime.now().timestamp()))
items.append(
ProductSummary(
goods_id=goods_id,
goods_name=goods_name,
goods_image=goods_image,
goods_price=goods_price,
third_shop_price=third_shop_price,
goods_link=goods_link,
on_sold=on_sold,
created_at=now,
updated_at=now,
)
)
return items
@staticmethod
def _extract_product_id(text: str) -> str | None:
"""从 URL 或文本中提取商品 ID
匹配 /product/detail/{id} 或 /product/other/{id} 格式
"""
match = re.search(r"/product/(?:detail|other)/([A-Za-z0-9]+)", text)
return match.group(1) if match else None
@staticmethod
def _normalize_url(url: str) -> str:
"""将相对 URL 补全为完整的 https URL"""
if url.startswith("http://") or url.startswith("https://"):
return url
return f"{BASE_URL}{url}"
def get_breadcrumb_list(self, tree):
dom_list = tree.css(".block-blocksurugayabreadcrum nav .breadcrumb .breadcrumb-item")
breadcrumb_list = []
for index, dom in enumerate(dom_list):
a_node = dom.css_first("a")
if not a_node:
continue
breadcrumb_list.append({
"position": index + 1,
"is_active": index == len(dom_list) - 1,
"name": a_node.text().strip(),
"href": a_node.attributes.get("href") or "",
})
return breadcrumb_list