import random import re import signal import socket import sys import time from decimal import Decimal, InvalidOperation from urllib.parse import quote from DrissionPage import ChromiumPage, ChromiumOptions import json import hashlib import os from base64 import b64decode from commons.Logger import get_spider_logger from commons.conn_mysql import MySQLPoolOn2 from pipelines.drug_pipelines import DrugPipeline from commons.feishu_webhook import send_text from spiders.jd.jd_captcha import handle_jd_slider_captcha from oss_upload.oss_upload import AliyunOSSUploader from commons.config import ( JD_DEVICE_ID, CHROME_PATH, UA_JD, CRAWLER_TOKEN, JD_NO_MATCH_THRESHOLD, JD_COLLECT_MAX_STEPS, ) import requests logger = get_spider_logger("jd") chrome_path = CHROME_PATH FETCH_TIMEOUT_FIRST = 3 FETCH_TIMEOUT_SCROLL = 3 LISTEN_CLEAR_ROUNDS = 3 LISTEN_CLEAR_TIMEOUT = 0.45 # 验证码 URL 特征 CAPTCHA_URL_KEYS = ["risk_handler", "cfe.m.jd.com", "verifycode", "captcha"] # 验证码页面文案特征 CAPTCHA_MARKERS = ["\u9a8c\u8bc1\u4e00\u4e0b\uff0c\u8d2d\u7269\u65e0\u5fe7", "\u524d\u65b9\u62e5\u6324", "\u62fc\u56fe\u9a8c\u8bc1"] # 「下一页」是否在视口内(条件略宽) _JS_NEXT_BTN_IN_VIEWPORT = """ var el = arguments[0]; if (!el) return false; var r = el.getBoundingClientRect(); var h = window.innerHeight || document.documentElement.clientHeight || 800; var w = window.innerWidth || document.documentElement.clientWidth || 1200; return r.bottom > 80 && r.top < h - 40 && r.right > 0 && r.left < w; """ class _PageInterrupted(Exception): """采集过程中页面出现异常(验证码/被踢/渲染异常),中断本页采集。""" def __init__(self, kind): self.kind = kind # "captcha" / "kicked" / "render" super().__init__(kind) class JdCrawlerV2: def __init__(self, drug_dict=None, scheduler=None, driver=None, cumulative_pages=0, cumulative_items=0, cumulative_stored=0, cumulative_skipped=0): self.driver = driver self.cumulative_pages = cumulative_pages self.cumulative_items = cumulative_items self.cumulative_stored = cumulative_stored self.cumulative_skipped = cumulative_skipped self.page_stored = 0 self.register_signal_handler() self.db = MySQLPoolOn2() self.ip = None self.account_name = None self.login_username = None self.login_password = None self.platform = 2 self.pipeline = DrugPipeline("jd") self.task_dict = drug_dict or {} self.scheduler = scheduler self.report_data = {} self.heartbeat_interval = 30 self.ossuploader = AliyunOSSUploader() self.start_page = 1 self.end_page = 1 if self.task_dict: self.get_product_data() self.success = True self.is_no_prodcut = 0 self._mouse = (500, 400) # 虚拟鼠标位置(视口坐标) self._snap_miss = 0 # 连续找不到商品元素的计数 self._snap_local_count = 1 # 本地快照序号 def get_product_data(self): self.task_id = self.task_dict["id"] self.company_id = self.task_dict["company_id"] self.product = self.task_dict["product_name"] self.product_desc = self.task_dict.get("product_specs", "") self.brand = self.task_dict.get("product_brand", "") self.product_keyword = self.task_dict.get("product_keyword", "") self.collect_task_id = self.task_dict.get("collect_task_id", "") self.sampling_cycle = self.task_dict.get("sampling_cycle", "") self.sampling_start_time = self.task_dict.get("sampling_start_time", "") self.sampling_end_time = self.task_dict.get("sampling_end_time", "") self.collect_equipment_id = self.task_dict.get("collect_equipment_id", "") self.account_id = self.task_dict.get("collect_equipment_account_id", "15") self.collect_region_id = self.task_dict.get("collect_region_id", "") self.collect_round = self.task_dict.get("collect_round", 1) self.start_page = self._parse_page(self.task_dict.get("start_page"), 1) self.end_page = 100 self.report_data = {'task_id': self.task_id, 'platform': self.platform, 'username': self.task_dict.get("username", JD_DEVICE_ID)} @staticmethod def _parse_page(value, default=1): try: page = int(value) return page if page >= 1 else default except (TypeError, ValueError): return default @staticmethod def _get_free_port(): """获取一个当前可用的本地端口,供 Chrome 调试使用。""" with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s: s.bind(("127.0.0.1", 0)) return s.getsockname()[1] def init_browser(self): if self.driver: logger.info("复用已有浏览器实例") self._listen_started = False return co = ChromiumOptions().set_browser_path(chrome_path) debug_port = self._get_free_port() co.set_user_data_path(f"./account_cache/jd/JD_101") co.set_local_port(debug_port) co.set_argument(f"--remote-debugging-port={debug_port}") co.set_argument("--remote-debugging-address=127.0.0.1") # co.set_argument("--disable-blink-features=AutomationControlled") co.set_argument("--disable-dev-shm-usage") co.set_argument("--no-first-run") # 避免首次运行弹窗 co.set_argument("--no-default-browser-check") # 避免默认浏览器检查 if self.ip: proxy = self.ip.strip() if not proxy.startswith(("http://", "https://")): proxy = f"http://{proxy}" co.set_argument(f"--proxy-server={proxy}") logger.info("启动浏览器: account=%s, debug_port=%s", self.account_name, debug_port) self.driver = ChromiumPage(co) self._listen_started = False def _start_listen(self): """登录完成后再开监听,避免干扰登录页/验证码拖动。""" if self._listen_started or not self.driver: return self.driver.listen.start("api?appid=search-pc-java") self._listen_started = True logger.info("已启动搜索接口监听") def register_signal_handler(self): def handler(signum, frame): print("\n⚠️ 程序退出") if self.driver: self.driver.quit() sys.exit(0) signal.signal(signal.SIGINT, handler) if hasattr(signal, "SIGTERM"): signal.signal(signal.SIGTERM, handler) def sleep(self, a, b): time.sleep(random.uniform(a, b)) def _scroll_page_down(self, delta=900): """拟人滚轮:真实 CDP mouseWheel 事件,物理签名。""" try: self._human_wheel(int(delta)) except Exception: self.driver.run_js(f"window.scrollBy(0, {int(delta)});") time.sleep(random.uniform(0.15, 0.35)) def _scroll_next_into_view(self, el): if not el: return try: self.driver.run_js( "arguments[0].scrollIntoView({block:'center',behavior:'instant'});", el, ) self.sleep(1, 2) except Exception as e: logger.warning("滚动到下一页按钮失败: %s", e) try: el.scroll.to_see() except Exception: pass def _get_scroll_info(self): return self.driver.run_js(""" return { scrollY: window.scrollY || window.pageYOffset || 0, docH: Math.max(document.body.scrollHeight, document.documentElement.scrollHeight, document.body.offsetHeight), viewH: window.innerHeight || document.documentElement.clientHeight || 800 }; """) def _find_next_btn(self, timeout=0.3): try: return self.driver.ele("text=下一页", timeout=timeout) except Exception: return None def _is_next_btn_visible(self, btn): if not btn: return False try: return bool(self.driver.run_js(_JS_NEXT_BTN_IN_VIEWPORT, btn)) except Exception: return False def _human_click(self, element): """拟人点击:贝塞尔移动+悬停+遮挡检测+真实 CDP mousePressed/Released。""" if not element: return False try: return self._human_click_el(element) except Exception as e: logger.warning("拟人点击失败,退回 JS click: %s", e) try: self.driver.run_js( "arguments[0].scrollIntoView({block:'center',behavior:'instant'});", element, ) self.driver.run_js("arguments[0].click();", element) return True except Exception: try: element.click() return True except Exception: return False @staticmethod def _estimated_price(json_data): fp = json_data.get("finalPrice") if isinstance(fp, dict): return fp.get("estimatedPrice", "") or "" return "" def get_heshu(self, full_title): last_box = None last_bottle = None for match in re.finditer(r"(\d+)(盒|瓶)", full_title): if match.group(2) == '盒': last_box = match else: # 瓶 last_bottle = match if last_box: return int(last_box.group(1)) elif last_bottle: return int(last_bottle.group(1)) else: return 1 def _take_snapshot(self, upload_key, ele): """在指定标签页截图并上传。""" # 拟人:滚轮把元素带进视口 → 鼠标移到商品上(像人在看) → 随机停顿 try: self._human_scroll_el_into_view(ele) pt = self._el_click_point(ele) if pt: self._human_move_to(*pt) if random.random() < 0.08: self.sleep(0.4, 0.9) # 偶尔停下来细看 else: self.sleep(0.05, 0.2) except Exception: pass time.sleep(0.05) # 快照前固定等待(过短可能截到未加载完的图片) # CDP clip 截图:避免 ele.get_screenshot 触发 Chrome 视口重排导致页面跳动 try: rect = self.driver.run_js( "var r=arguments[0].getBoundingClientRect();" "return [r.left + window.scrollX, r.top + window.scrollY, r.width, r.height," "r.top, r.bottom, r.left, r.right, window.innerHeight, window.innerWidth, window.scrollY];", ele, ) if not rect or rect[2] <= 0 or rect[3] <= 0: logger.warning("元素区域异常 upload_key=%s", upload_key) return "" x, y, w, h = rect[0], rect[1], rect[2], rect[3] clip = {"x": x, "y": y, "width": w, "height": h, "scale": 1} # 元素完整在视口内(垂直+水平)才不开 captureBeyondViewport in_view = (rect[4] >= 0 and rect[5] <= rect[8] and rect[6] >= 0 and rect[7] <= rect[9]) if in_view: data = self.driver.run_cdp( "Page.captureScreenshot", format="jpeg", quality=100, clip=clip, _timeout=12, )["data"] else: data = self.driver.run_cdp( "Page.captureScreenshot", format="jpeg", quality=100, captureBeyondViewport=True, clip=clip, _timeout=12, )["data"] self.driver.run_js( "window.scrollTo({top: %d, behavior: 'instant'});" % int(rect[10]) ) jpg_bytes = b64decode(data) if not jpg_bytes: logger.warning("截图为空 upload_key=%s", upload_key) return "" self._last_snap_bytes = jpg_bytes # 暂存,供 parse 内本地保存 img_url = self.ossuploader.upload_from_bytes(jpg_bytes, str(upload_key)) except Exception: logger.exception("截图或 OSS 上传失败 upload_key=%s", upload_key) return "" if not img_url: logger.warning("OSS 未返回有效地址 upload_key=%s", upload_key) return "" logger.info("截图上传完成 upload_key=%s url=%s", upload_key, img_url) time.sleep(random.uniform(0.05, 0.15)) return img_url def parse(self, ware_list, force=False): unprocessed = [] for w in ware_list: # 每5条做一次页面级异常巡检,命中立刻中断本页采集 if len(unprocessed) % 5 == 0: intr = self._interruption() if intr: raise _PageInterrupted(intr) if not isinstance(w, dict): continue sku_id = str(w.get("skuId", "")) if not sku_id: continue ele_xpath = "//div[@id='main_search_conter']//div[contains(@class,'_goodsContainer_')]/div[@data-sku=" + "'" + sku_id + "'" + "]" ele_screen = self.driver.ele("xpath=" + ele_xpath, timeout=0.2) if not ele_screen and not force: unprocessed.append(w) continue title = w.get("wareName", "") title = re.sub(r"<[^>]*>", "", title).strip() color = w.get("color", "") full_title = title + " " + color logger.info(full_title) if self.product not in full_title: self.is_no_prodcut += 1 continue if self.brand not in full_title: self.is_no_prodcut += 1 continue if self.product_desc: if self.product_desc in full_title: crawl_product_desc = self.product_desc else: crawl_product_desc = "" title = full_title else: crawl_product_desc = "" title = full_title if "+[" in title: continue self.is_no_prodcut = 0 status = 1 if self.product_keyword: search_keyword_list = self.product_keyword.split(",") for search_keyword in search_keyword_list: if search_keyword.strip() not in title: status = 0 if status == 0: continue logger.info(f"商品名:{title}") sku_id = w.get("skuId", "") sales = w.get("totalSales", "") shop_id = w.get("shopId", "") shop_name = w.get("shopName", "") heshu_count = self.get_heshu(full_title) final_price = self._estimated_price(w) jd_price = w.get("jdPrice", "") item_url = f"https://item.jd.com/{sku_id}.html" low_price = final_price if final_price else jd_price # 获取列表页快照 upload_key = hashlib.md5(item_url.encode("utf-8")).hexdigest() snap_url = "" if ele_screen: for i in range(3): snap_url = self._take_snapshot(upload_key, ele_screen) if snap_url != "": break else: logger.warning(f"未找到商品元素无法截图: {sku_id}") try: price = Decimal(str(low_price)).quantize(Decimal("0.00")) except (InvalidOperation, ValueError): price = Decimal("0.00") # 本地留存快照用于人工检查质量: img/001_价格.jpg if snap_url and getattr(self, '_last_snap_bytes', None): try: local_dir = os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "..", "img") os.makedirs(local_dir, exist_ok=True) local_name = f"{self._snap_local_count:03d}_{price}.jpg" with open(os.path.join(local_dir, local_name), "wb") as lf: lf.write(self._last_snap_bytes) self._snap_local_count += 1 self._last_snap_bytes = None except Exception: pass item_url = f"https://item.jd.com/{sku_id}.html" mall_url = f"https://mall.jd.com/index-{shop_id}.html?from=pc" # 字段与 yaofangwang_crawl 对齐;键顺序须与 commons.sql_data.RETRIEVE_SCRAPE_INSERT_COLUMNS 一致 now_ts = time.strftime("%Y-%m-%d %H:%M:%S") product = { "platform": self.platform, "item_id": sku_id, "enterprise_id": self.company_id, "product_name": title, "spec": crawl_product_desc, "one_price": "", "detail_url": item_url, "shop_name": shop_name, "anonymous_store_name": "", "shop_url": mall_url, "city_name": "", "city_id": "", "province_name": "", "province_id": "", "shipment_city_name": "", "shipment_city_id": "", "shipment_province_name": "", "shipment_province_id": "", "area_info": "", "factory_name": "", "scrape_date": time.strftime("%Y-%m-%d"), "price": price, "sales": sales, "stock_count": "", "snapshot_url": snap_url, "approval_num": "", "produced_time": "", "deadline": "", "update_time": now_ts, "insert_time": now_ts, "number": heshu_count, "product_brand": self.brand or "", "collect_task_id": self.collect_task_id, "task_id": self.task_id, "search_name": self.product, "company_name": "", "collect_config_info": json.dumps( { "sampling_cycle": self.sampling_cycle, "sampling_start_time": self.sampling_start_time, "sampling_end_time": self.sampling_end_time, } ), "account_id": self.account_id, "collect_region_id": self.collect_region_id, "collect_round": self.collect_round, "is_sold_out": 0 } try: affected_rows = self.pipeline.storge_data(product) if affected_rows and affected_rows > 0: self.page_stored += 1 logger.info("%s", json.dumps(product, ensure_ascii=False, default=str)) except Exception as e: logger.exception("写入数据库失败: %s", e) return unprocessed @staticmethod def _response_has_ware_list(data): if not isinstance(data, dict): return False inner_data = data.get("data") if not isinstance(inner_data, dict): return False return bool(inner_data.get("wareList")) def fetch_items_once(self, timeout=FETCH_TIMEOUT_FIRST): wares = [] for resp in self.driver.listen.steps(timeout=timeout): try: data = resp.response.body if not self._response_has_ware_list(data): continue ware_list = data["data"]["wareList"] wares.extend(ware_list) except Exception as e: logger.warning("解析监听响应失败: %s", e) return wares def clear_listen_buffer(self, rounds=LISTEN_CLEAR_ROUNDS, timeout=LISTEN_CLEAR_TIMEOUT): try: for _ in range(rounds): resps = list(self.driver.listen.steps(timeout=timeout)) if not resps: break logger.debug("监听缓冲已清空") except Exception as e: logger.debug("清空监听缓冲失败: %s", e) def collect_full_page_items(self, max_steps=JD_COLLECT_MAX_STEPS): """单次循环:边滑动边收数据,到底 / 看见「下一页」即停。""" # 页首巡检 + 清理误开标签页 self._close_unexpected_tabs() intr = self._interruption() if intr: raise _PageInterrupted(intr) pending_wares = [] total_n = 0 new_wares = self.fetch_items_once(timeout=FETCH_TIMEOUT_FIRST) total_n += len(new_wares) pending_wares.extend(new_wares) pending_wares = self.parse(pending_wares) stagnant = 0 last_scroll_y = None for step in range(max_steps): # 每步巡检:中途出验证码/被踢立刻中断本页 intr = self._interruption() if intr: raise _PageInterrupted(intr) next_btn = self._find_next_btn(timeout=0.3) if self._is_next_btn_visible(next_btn): new_wares = self.fetch_items_once(timeout=FETCH_TIMEOUT_SCROLL) total_n += len(new_wares) pending_wares.extend(new_wares) pending_wares = self.parse(pending_wares) self.parse(pending_wares, force=True) return total_n, next_btn info = self._get_scroll_info() scroll_y = info["scrollY"] doc_h = info["docH"] view_h = info["viewH"] at_bottom = (scroll_y + view_h >= doc_h - 20) if last_scroll_y is not None and abs(scroll_y - last_scroll_y) < 8: stagnant += 1 else: stagnant = 0 last_scroll_y = scroll_y if at_bottom and stagnant >= 2: new_wares = self.fetch_items_once(timeout=FETCH_TIMEOUT_SCROLL) total_n += len(new_wares) pending_wares.extend(new_wares) pending_wares = self.parse(pending_wares) self.parse(pending_wares, force=True) next_btn = self._find_next_btn(timeout=2) if next_btn: self._scroll_next_into_view(next_btn) return total_n, next_btn logger.info("已到页面底部且未发现下一页,停止滑动") return total_n, None # 滚动前约1/4概率先把鼠标移近滚动区域(真人边读边滚,鼠标不是死的) if random.random() < 0.25: try: vw, vh = self.driver.run_js("return [window.innerWidth, window.innerHeight];") self._human_move_to( vw * random.uniform(0.25, 0.75), vh * random.uniform(0.3, 0.7) ) except Exception: pass self._scroll_page_down(random.randint(600, 1100)) if random.random() < 0.15: self._human_wheel(-random.randint(60, 140)) self.sleep(0.1, 0.35) pending_wares = self.parse(pending_wares) if step % 3 == 2: new_wares = self.fetch_items_once(timeout=FETCH_TIMEOUT_SCROLL) total_n += len(new_wares) pending_wares.extend(new_wares) pending_wares = self.parse(pending_wares) new_wares = self.fetch_items_once(timeout=FETCH_TIMEOUT_SCROLL) total_n += len(new_wares) pending_wares.extend(new_wares) pending_wares = self.parse(pending_wares) self.parse(pending_wares, force=True) next_btn = self._find_next_btn(timeout=3) if next_btn and not self._is_next_btn_visible(next_btn): self._scroll_next_into_view(next_btn) return total_n, next_btn def get_account(self): sql_account = """ SELECT * FROM `retrieve_collect_equipment_account` WHERE `id` = %s and `status` = 0 """ account_list = self.db.select_data(sql_account, self.account_id) if not account_list: return False account_dict = account_list[0] print(account_dict) self.ip = account_dict.get("ip") self.account_name = account_dict.get("username") self.login_username = account_dict.get("phone") or account_dict.get("username", "") self.login_password = account_dict.get("password", "") logger.info("获取到账号: %s, ip: %s", self.account_name, self.ip) return True def disable_account(self): update_sql = f""" UPDATE `retrieve_collect_equipment_account` SET `status`= %s WHERE `name` = %s; """ self.db.execute(update_sql, (1, self.account_name)) def _build_search_keyword(self): parts = [p for p in (self.brand, self.product, self.product_desc) if p] return " ".join(parts).strip() or self.product def _is_logged_out(self): return bool(self.driver.ele("xpath=//*[@class='link-login']", timeout=2)) # ==================== 拟人化:鼠标层 ==================== def _human_move_to(self, x, y, duration=None): """贝塞尔曲线移动鼠标(ease-out),沿途产生真实 CDP mousemove 事件。""" x0, y0 = self._mouse dist = ((x - x0) ** 2 + (y - y0) ** 2) ** 0.5 if dist < 2: self._mouse = (x, y) return c1 = ( x0 + (x - x0) * random.uniform(0.15, 0.4) + random.uniform(-0.15, 0.15) * dist, y0 + (y - y0) * random.uniform(0.15, 0.4) + random.uniform(-0.15, 0.15) * dist, ) c2 = ( x0 + (x - x0) * random.uniform(0.6, 0.9) + random.uniform(-0.15, 0.15) * dist, y0 + (y - y0) * random.uniform(0.6, 0.9) + random.uniform(-0.15, 0.15) * dist, ) steps = max(10, min(35, int(dist / 25))) if duration is None: duration = max(0.08, min(0.45, dist / random.uniform(2200, 3800))) for i in range(1, steps + 1): t = 1 - (1 - i / steps) ** 2 mt = 1 - t px = mt ** 3 * x0 + 3 * mt ** 2 * t * c1[0] + 3 * mt * t ** 2 * c2[0] + t ** 3 * x py = mt ** 3 * y0 + 3 * mt ** 2 * t * c1[1] + 3 * mt * t ** 2 * c2[1] + t ** 3 * y self.driver.run_cdp("Input.dispatchMouseEvent", type="mouseMoved", x=px, y=py) time.sleep(duration / steps * random.uniform(0.6, 1.4)) self._mouse = (x, y) def _el_click_point(self, el): """元素可视区域内的随机点击点(视口坐标),不在视口返回 None。""" try: r = self.driver.run_js( "var r=arguments[0].getBoundingClientRect();" "return [r.left, r.top, r.width, r.height, window.innerWidth, window.innerHeight];", el, ) except Exception: return None if not r: return None left, top, w, h, vw, vh = r x = left + w * random.uniform(0.25, 0.75) y = top + h * random.uniform(0.3, 0.7) if not (2 <= x <= vw - 2 and 2 <= y <= vh - 2): return None return (x, y) def _hit_check(self, el, x, y): """确认 (x,y) 处顶层元素是目标本身/子元素;否则返回遮挡物描述。""" try: return self.driver.run_js( "var t=arguments[2];" "var h=document.elementFromPoint(arguments[0],arguments[1]);" "if(!h) return 'none';" "if(h===t||t.contains(h)||h.contains(t)) return 'ok';" "var d=(h.tagName||'').toLowerCase();" "if(h.id) d+='#'+h.id;" "if(h.className&&typeof h.className==='string') d+='.'+h.className.trim().split(' ')[0];" "return d;", x, y, el, ) except Exception: return "err" def _human_click_el(self, el): """真人点击:自然滚到元素→贝塞尔移动→悬停→遮挡检测→真实 mousePressed/Released。""" for attempt in range(3): if not self._human_scroll_el_into_view(el): try: self.driver.run_js( "arguments[0].scrollIntoView({block:'center',behavior:'instant'});", el ) except Exception: pass self.sleep(0.15, 0.35) self._wait_scroll_settle() pt = self._el_click_point(el) if not pt: return False self._human_move_to(*pt) self.sleep(0.2, 0.6) hit = self._hit_check(el, *pt) if hit == "ok": x, y = pt self.driver.run_cdp( "Input.dispatchMouseEvent", type="mousePressed", x=x, y=y, button="left", clickCount=1, ) time.sleep(random.uniform(0.05, 0.12)) self.driver.run_cdp( "Input.dispatchMouseEvent", type="mouseReleased", x=x, y=y, button="left", clickCount=1, ) self._close_unexpected_tabs() return True logger.info("点击被遮挡(%s),按 Esc 关弹层重试(%d/3)", hit, attempt + 1) self._press_escape() self.sleep(0.5, 1.2) if attempt >= 1: self._hide_floater_at(*pt, el) logger.warning("多次重试仍被遮挡,放弃本次点击") return False def _mouse_wander(self): """采集间隙随机移动鼠标,模拟真人浏览时的手部小动作。""" try: loc = self.driver.run_js( "return [Math.round(window.scrollX + innerWidth * Math.random())," + " Math.round(window.scrollY + innerHeight * Math.random())];" ) self.driver.actions.move_to(loc, duration=random.uniform(0.4, 1.0)) except Exception: pass # ==================== 拟人化:滚轮层 ==================== def _human_wheel(self, delta): """真实滚轮事件:120/格物理签名,1~4格一串,格间隔30~70ms。""" try: vw, vh = self.driver.run_js("return [window.innerWidth, window.innerHeight];") except Exception: vw, vh = 1200, 800 x, y = self._mouse if not (0 <= x <= vw and 0 <= y <= vh): self._human_move_to(vw * random.uniform(0.3, 0.7), vh * random.uniform(0.3, 0.6)) x, y = self._mouse sign = 1 if delta >= 0 else -1 notches = max(1, round(abs(delta) / 120)) while notches > 0: burst = min(notches, random.randint(1, 4)) for _ in range(burst): self.driver.run_cdp( "Input.dispatchMouseEvent", type="mouseWheel", x=x, y=y, deltaX=0, deltaY=sign * 120, ) time.sleep(random.uniform(0.03, 0.07)) notches -= burst if notches > 0: time.sleep(random.uniform(0.04, 0.10)) time.sleep(random.uniform(0.05, 0.12)) def _human_scroll_el_into_view(self, el, max_tries=6): """用自然滚轮把元素带进视口(不用 scrollIntoView 瞬移)。""" for _ in range(max_tries): r = self.driver.run_js( "var r=arguments[0].getBoundingClientRect();" "return [r.top, r.bottom, r.left, r.right, window.innerHeight, window.innerWidth];", el, ) if not r: return False top, bottom, left, right, vh, vw = r # ---- 垂直方向 ---- el_h = bottom - top if el_h > vh - 60: # 高元素:视口放不下,顶部对齐即可 v_ok = (8 <= top <= 120) v_delta = 0 if v_ok else int(top - 40) else: # 完整可见才算到位 v_ok = (top >= 8 and bottom <= vh - 8) v_delta = 0 if v_ok else int((top + bottom) / 2 - vh / 2) # ---- 水平方向 ---- el_w = right - left if el_w > vw - 16: # 宽元素:视口放不下,左边缘对齐即可 h_ok = (0 <= left <= 30) h_delta = 0 if h_ok else int(left - 10) else: # 完整可见才算到位 h_ok = (left >= 0 and right <= vw - 4) h_delta = 0 if h_ok else int((left + right) / 2 - vw / 2) if v_ok and h_ok: return True # 先调垂直(自然滚轮),再调水平(window.scrollBy:左右滚动没有滚轮物理事件) if v_delta != 0: self._human_wheel(max(-600, min(600, v_delta))) if h_delta != 0: # 水平滚动:使用 deltaX(触控板/高端鼠标本身就有水平滚轮,不做 Shift 修饰) sign = 1 if h_delta >= 0 else -1 notches = max(1, round(abs(h_delta) / 120)) while notches > 0: burst = min(notches, random.randint(1, 4)) for _ in range(burst): self.driver.run_cdp( "Input.dispatchMouseEvent", type="mouseWheel", x=self._mouse[0], y=self._mouse[1], deltaX=sign * 120, deltaY=0, ) time.sleep(random.uniform(0.03, 0.07)) notches -= burst if notches > 0: time.sleep(random.uniform(0.04, 0.10)) time.sleep(random.uniform(0.05, 0.12)) self._wait_scroll_settle() # 等页面滚动惯性完全停下再读坐标,避免振荡 return False def _wait_scroll_settle(self, timeout=3): """等滚动惯性停下再取坐标。""" t0 = time.time() last = None while time.time() - t0 < timeout: try: y = self.driver.run_js("return Math.round(window.scrollY);") except Exception: return if last is not None and abs(y - last) < 2: return last = y time.sleep(0.15) # ==================== 拟人化:键盘/浮层/标签页 ==================== def _human_type_digits(self, el, text): """拟人输入数字:点击聚焦→Ctrl+A全选→逐键 CDP rawKeyDown+char+keyUp。""" if not self._human_click_el(el): el.click() time.sleep(random.uniform(0.15, 0.4)) self.driver.run_cdp( "Input.dispatchKeyEvent", type="rawKeyDown", key="a", code="KeyA", modifiers=2, windowsVirtualKeyCode=65, ) self.driver.run_cdp( "Input.dispatchKeyEvent", type="keyUp", key="a", code="KeyA", modifiers=2, windowsVirtualKeyCode=65, ) time.sleep(random.uniform(0.05, 0.15)) for ch in text: vk = ord(ch) code = f"Digit{ch}" self.driver.run_cdp( "Input.dispatchKeyEvent", type="rawKeyDown", key=ch, code=code, windowsVirtualKeyCode=vk, ) self.driver.run_cdp( "Input.dispatchKeyEvent", type="char", text=ch, key=ch, code=code, windowsVirtualKeyCode=vk, ) self.driver.run_cdp( "Input.dispatchKeyEvent", type="keyUp", key=ch, code=code, windowsVirtualKeyCode=vk, ) time.sleep(random.uniform(0.06, 0.18)) def _search_via_box(self, keyword, kw): """拟人搜索路径:搜索框点击→全选清空→insertText注入中文→回车提交。""" search_url = f"https://search.jd.com/Search?keyword={kw}&enc=utf-8&wq={kw}" self._close_unexpected_tabs() url_now = self.driver.url or "" if "www.jd.com" not in url_now and "search.jd.com" not in url_now: self.driver.get("https://www.jd.com/", timeout=15) self.sleep(3, 5) box = None for loc in ("css=#key", "css=input[placeholder*='\u641c\u7d22']", "css=.search-m input[type='text']", "css=input[name='keyword']"): box = self.driver.ele(loc, timeout=2) if box: break if not box: box = self.driver.run_js( "var ins=document.querySelectorAll(\"input[type='text'],input:not([type])\");" "for(var i=0;i100&&r.height>15&&r.top>=-50&&r.top 20: last_tip = time.time() logger.info("[验证码] 等待人工验证中…… 已等待 %ds", int(time.time() - t0)) time.sleep(3) logger.warning("[验证码] 等待人工验证超时") return False def _recover_from_captcha(self, page_no, kw): """验证码恢复:等人工验证通过后跳回原页续采。""" result = self._wait_captcha_solved() if result == "kicked": return self._recover_from_kick(page_no, kw) if not result: return False return self._reposition_to_page(page_no, kw) def _recover_from_kick(self, page_no, kw): """被踢恢复:先试自动登录,失败则等人工扫码,成功后跳回原页。""" logger.info("[被踢] 第 %d 页检测到被踢出登录", page_no) # 先试自动登录 if self.login_password and self.login_username: if self._ensure_logged_in(): logger.info("[被踢] 自动登录成功,跳回第 %d 页", page_no) return self._reposition_to_page(page_no, kw) # 自动登录失败/无凭据:等人工扫码 logger.info("[被踢] 请在浏览器中扫码登录,脚本将持续等待……") t0 = time.time() while time.time() - t0 < 3600: url = self.driver.url or "" if "passport.jd.com" not in url: self.sleep(2, 4) if not self._is_logged_out(): logger.info("[被踢] 人工登录成功,跳回第 %d 页", page_no) return self._reposition_to_page(page_no, kw) time.sleep(5) logger.warning("[被踢] 等待人工登录超时") return False def _recover_from_empty(self, page_no, kw): """空页恢复:像人一样刷新重试,访问频繁时切长间隔。""" logger.info("[空页] 第 %d 页搜索返回空结果,像人一样刷新重试", page_no) for attempt in range(1, 4): if "\u8bbf\u95ee\u9891\u7e41" in (self.driver.html or ""): wait = {1: 90, 2: 180, 3: 300}.get(attempt, 300) logger.info("[空页] 检测到访问频繁,等 %ds 后重试", wait) time.sleep(wait) else: self.sleep(18, 35) self.driver.refresh() self.sleep(6, 10) if not self._is_empty_result() and not self._has_captcha(): logger.info("[空页] 第 %d 次刷新后数据恢复", attempt) return True logger.info("[空页] 第 %d 次刷新后仍为空", attempt) logger.warning("[空页] 多次刷新仍无数据,停止(今日额度可能已尽)") return False def _error_report(self, data): self.report_data.update(data) if self.scheduler: self.scheduler.stop() self.scheduler.post_report(self.report_data) def _success_report(self, data): self.report_data.update(data) if self.scheduler: self.scheduler.post_report(self.report_data) def post_report(self, data): url = "http://192.168.2.246:8080/api/collect_task/report" print('传给返回接口的数据', data) headers = {'X-Crawler-Token': CRAWLER_TOKEN} response = requests.post(url, json=data, headers=headers) result = response.json() if result.get('code') != 'success': logger.info(f'翻页回传:{result}') if self.scheduler: self.scheduler.stop() print(f'任务进度上传 {result}') def perform_jd_login(self): """ 使用已有浏览器实例执行京东账号密码登录(含滑块验证码)。 成功返回 True,失败返回 False。 """ username = self.login_username password = self.login_password login_url = "https://passport.jd.com/new/login.aspx" self.driver.get(login_url) input_name = self.driver.ele("xpath=//input[@id='loginname']", timeout=15) if not input_name: print("未找到用户名输入框") return False input_name.input(username) time.sleep(random.uniform(1.5, 2.5)) input_pass = self.driver.ele("xpath://input[@name='nloginpwd']", timeout=5) if not input_pass: print("未找到密码输入框") return False input_pass.input(password) time.sleep(random.uniform(1.5, 2.5)) login_btn = self.driver.ele("xpath://a[@id='loginsubmit']", timeout=5) if not login_btn: print("未找到登录按钮") return False login_btn.click() time.sleep(random.uniform(3, 5)) if not handle_jd_slider_captcha(self.driver): print("滑块验证码未通过") return False return True def _ensure_logged_in(self): """未登录时自动走登录流程(账号密码 + 滑块)。""" if not self._is_logged_out(): return True logger.info("检测到未登录,开始自动登录: %s", self.account_name) ok = self.perform_jd_login() if ok and not self._is_logged_out(): logger.info("自动登录成功: %s", self.account_name) return True logger.error("自动登录失败: %s", self.account_name) return False def _check_page_blocked(self): html = self.driver.html or "" if "抱歉由于访问频繁导致无法搜索" in html: logger.error("账号无法搜索(访问频繁)") self.success = False return True return False def _jump_to_page(self, target_page): """跳转到指定页码,并清空跳转前的监听残留。""" to_page_input = self.driver.ele( "xpath=//input[contains(@class, 'pagination-input')] | //div[contains(@class,'_pagination_toPageNum_')]//input[@type='text']", timeout=3, ) if not to_page_input: logger.warning("未找到跳页输入框,无法跳转到第 %s 页", target_page) return False self.clear_listen_buffer() # 拟人输入:逐键 CDP 打字(失败自动退回直接赋值) try: self._human_type_digits(to_page_input, str(target_page)) except Exception: to_page_input.clear() to_page_input.input(str(target_page)) self.sleep(0.4, 0.8) confirm_btn = self.driver.ele("xpath=//button[contains(@class, 'pagination-jump-btn')]", timeout=1) if confirm_btn: self.driver.run_js("arguments[0].click();", confirm_btn) else: self.driver.actions.key_down("enter").key_up("enter") self.sleep(1.5, 2.5) self.clear_listen_buffer() logger.info("已跳转到第 %s 页", target_page) return True def _go_next_page(self, next_btn): self.clear_listen_buffer() if not self._human_click(next_btn): logger.warning("点击下一页失败") return False self._close_unexpected_tabs() self.sleep(0.6, 1.2) return True def crawl(self): total = 0 keyword = self._build_search_keyword() # 如果已在京东页面(搜索页/首页)直接用当前页的搜索框,不必每次都回首页 url_now = (self.driver.url or "").lower() if "jd.com" not in url_now: self.driver.get("https://www.jd.com/", timeout=15) self.sleep(3, 5) if self._is_logged_out(): # 打开扫码登录页,等人工扫码(实时轮询,登录即继续,最多等5分钟) logger.info("未登录,打开扫码登录页,请用京东APP扫码……") self.driver.get("https://passport.jd.com/new/login.aspx", timeout=15) t0 = time.time() logged_in = False while time.time() - t0 < 300: # 最多等5分钟 if "passport.jd.com" not in (self.driver.url or ""): self.sleep(2, 4) if not self._is_logged_out(): logged_in = True break time.sleep(3) if logged_in: logger.info("人工登录成功") else: logger.error("等待人工登录超时(5分钟)") if self.login_password and self.login_username: self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': self.start_page, 'exception_type': 1}) else: self.disable_account() send_text(f"京东:{self.account_name}账号登录失败") self.success = False return kw = quote(str(keyword or ""), safe="") self._search_kw = kw # 必须先监听再打开搜索页,否则首屏 wareList(前约 30 条)在监听开启前就返回了 self._start_listen() # 拟人搜索:搜索框输入(失败自动退回 URL 直达) self._search_via_box(keyword, kw) self.sleep(2, 3) if self._check_page_blocked(): self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': self.start_page, 'exception_type': 2}) return if not handle_jd_slider_captcha(self.driver, pause_listen=False): logger.warning("进入搜索页后滑块验证码处理失败") self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': self.start_page, 'exception_type': 2}) self.success = False return if self.start_page > 1: if not self._jump_to_page(self.start_page): logger.warning("跳页失败,将从第 1 页开始采集") self.start_page = 1 logger.info( "采集页码范围: %s ~ %s(共 %s 页)", self.start_page, self.end_page, self.end_page - self.start_page + 1, ) page_no = self.start_page while page_no <= self.end_page: # ---- 页首异常检测 + 恢复 ---- intr = self._interruption() if intr == "kicked": if self._recover_from_kick(page_no, kw): continue self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 1}) break elif intr == "captcha": if self._recover_from_captcha(page_no, kw): continue self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 2}) break elif self._is_empty_result(): if self._recover_from_empty(page_no, kw): continue self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 2}) break if self._is_logged_out(): if not self._ensure_logged_in(): self.success = False break self._search_via_box(keyword, kw) self.sleep(3, 5) if page_no > 1: self._jump_to_page(page_no) if not handle_jd_slider_captcha(self.driver, pause_listen=True): logger.warning("滑块验证码处理失败,停止采集") self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 2}) self.success = False break if self._check_page_blocked(): break logger.info("===== 正在爬取第 %s 页 =====", page_no) self.page_stored = 0 search_ele = self.driver.ele("xpath=//div[@id='search-condition']", timeout=10) if not search_ele: logger.warning("未找到搜索结果区域,停止采集") self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 2}) break # ---- 采集 ---- try: page_n, _ = self.collect_full_page_items() except _PageInterrupted as pi: if pi.kind == "captcha": if not self._recover_from_captcha(page_no, kw): self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 2}) break continue if pi.kind == "render": logger.warning("第 %d 页渲染异常,重新加载本页重采", page_no) if not self._reposition_to_page(page_no, kw): break continue # kicked if not self._recover_from_kick(page_no, kw): self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': page_no, 'exception_type': 1}) break continue logger.info("本页监听商品条数(含可能重复): %s", page_n) total += page_n logger.info("累计监听条数: %s", total) if self.scheduler and self.scheduler.end: logger.info('心跳失败') self._error_report({'is_finished': 0, 'need_reassign': 0, 'current_page': page_no}) break if self.is_no_prodcut > JD_NO_MATCH_THRESHOLD: logger.info("连续无匹配商品过多,停止采集") self._success_report({'is_finished': 1, 'need_reassign': 0, 'current_page': page_no}) break self.cumulative_pages += 1 self.cumulative_items += page_n page_skipped = page_n - self.page_stored self.cumulative_stored += self.page_stored self.cumulative_skipped += page_skipped logger.info( "关键字 %s 第 %s 页获取完成, 本页获取量: %d, 本页已入库: %d, 本页未入库: %d | 账号测试总页数: %d, 账号测试总数据量: %d, 总已入库: %d, 总未入库: %d", keyword, page_no, page_n, self.page_stored, page_skipped, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped ) print(f"[{keyword}] 当前页数: {page_no}, 本页获取数据量: {page_n}, 本页已入库: {self.page_stored}, 本页未入库: {page_skipped} | 账号测试总页数: {self.cumulative_pages}, 账号测试总数据量: {self.cumulative_items}, 总已入库: {self.cumulative_stored}, 总未入库: {self.cumulative_skipped}") if page_no >= self.end_page: self._success_report({'is_finished': 1, 'need_reassign': 0, 'current_page': page_no}) break next_btn = self.driver.ele("text=下一页", timeout=2) if not next_btn: logger.info("没有下一页(未找到)") self._success_report({'is_finished': 1, 'need_reassign': 0, 'current_page': page_no}) break cls_str = next_btn.attr("class") or "" if "disabled" in cls_str: logger.info("没有下一页(已禁用)") self._success_report({'is_finished': 1, 'need_reassign': 0, 'current_page': page_no}) break if not self._go_next_page(next_btn): break page_no += 1 if self.scheduler: self.scheduler.stop() def run(self): # 检测账号 if not self.get_account(): logger.info("==================当前无账号可用==================") self.success = False return self.pipeline.crawl_count, self.success, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped logger.info("获取到账号:%s,代理ip:%s", self.account_name, self.ip) # # # 每次选取账号,立马账号使用时间 update_sql = f""" UPDATE `retrieve_collect_equipment_account` SET `status`= %s, `update_time`= %s WHERE `username` = %s; """ self.db.execute(update_sql, (0, int(time.time()), self.account_name)) try: self.init_browser() self.crawl() except Exception as e: self.success = False logger.exception("爬取异常: %s", e) self._error_report({'is_finished': 0, 'need_reassign': 1, 'current_page': self.start_page, 'exception_type': 3}) self.sleep(3, 5) finally: pass # if self.driver: # self.driver.quit() # self.driver = None return self.pipeline.crawl_count, self.success, self.driver, self.cumulative_pages, self.cumulative_items, self.cumulative_stored, self.cumulative_skipped