如何在Python多进程中调用Product类的parse方法?
多进程批量解析URL实现方案
核心思路
利用Python的multiprocessing模块创建进程池,并行处理多个URL的解析任务,同时修正原代码中类变量的潜在问题,确保多进程环境下实例完全隔离。
步骤1:修正Product类(关键)
原代码中soup和url是类变量,会被所有实例共享,在多进程环境下可能引发意外冲突,需改为实例变量:
class Product: def __init__(self, url, user_agents): self.url = url # 改为实例变量 self.soup = None # 改为实例变量 print('Class Initiated with URL: {}'.format(url)) # 随机化User-Agent user_agent = get_random_user_agent(user_agents) user_agent = user_agent.rstrip('\n') if 'linux' in user_agent.lower(): sec_ch_ua_platform = 'Linux' elif 'mac os x' in user_agent.lower(): sec_ch_ua_platform = 'macOS' else: sec_ch_ua_platform = 'Windows' headers = { 'User-Agent': user_agent, 'Sec-CH-UA-Platform': sec_ch_ua_platform # 补充其他需要的请求头字段 } r = create_request(url, None, headers=headers, is_proxy=False) if r is None: raise ValueError('Could not get data') html = r.text.strip() self.soup = BeautifulSoup(html, 'lxml') def parse(self): record = {} name = '' price = 0 user_count_in_cart = 0 review_count = 0 rating = 0 is_personalized = 'no' try: name = self.get_name() price = self.get_price() is_pick = self.get_is_pick() # 补充剩余字段的解析逻辑 record['name'] = name record['price'] = price record['is_pick'] = is_pick # ...其他字段 return record except Exception as e: print(f"解析URL {self.url} 失败: {str(e)}") return None # 补充实现get_name、get_price等方法(示例) def get_name(self): return self.soup.find('h1', class_='product-name').text.strip() def get_price(self): price_str = self.soup.find('span', class_='product-price').text.strip() return float(price_str.replace('$', '')) def get_is_pick(self): return 'yes' if self.soup.find('span', class_='pick-tag') else 'no'
步骤2:多进程实现(推荐用进程池)
使用multiprocessing.Pool可以快速实现批量并行处理,代码简洁易维护:
import multiprocessing # 封装单个URL的处理逻辑,适配进程池调用 def process_single_url(url, user_agents): try: product = Product(url, user_agents) return product.parse() except Exception as e: print(f"处理URL {url} 出错: {str(e)}") return None if __name__ == "__main__": # 假设你的URL列表和User-Agent列表 links = ["https://example.com/product1", "https://example.com/product2", "https://example.com/product3"] user_agents = [ "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Safari/537.36", "Mozilla/5.0 (Macintosh; Intel Mac OS X 13_4) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/16.5 Safari/605.1.15", "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Safari/537.36" ] # 设置进程数:建议等于CPU核心数,或根据目标网站反爬策略调整(避免请求过于频繁) process_count = multiprocessing.cpu_count() # 创建进程池并执行任务 with multiprocessing.Pool(process_count) as pool: # 使用starmap传递多参数(URL和User-Agent列表) results = pool.starmap(process_single_url, [(url, user_agents) for url in links]) # 过滤无效结果,收集有效数据 valid_records = [record for record in results if record is not None] print(f"成功解析 {len(valid_records)} 条商品数据") # 可将结果保存到文件或数据库 # import json # with open('products.json', 'w', encoding='utf-8') as f: # json.dump(valid_records, f, ensure_ascii=False, indent=2)
步骤3:手动管理进程(精细控制场景)
如果需要更精细的进程控制(如自定义进程名、资源限制),可直接使用multiprocessing.Process:
import multiprocessing def worker(url, user_agents, result_queue): try: product = Product(url, user_agents) result_queue.put(product.parse()) except Exception as e: print(f"处理URL {url} 出错: {str(e)}") result_queue.put(None) if __name__ == "__main__": links = ["https://example.com/product1", "https://example.com/product2"] user_agents = ["ua1", "ua2"] result_queue = multiprocessing.Queue() processes = [] # 创建并启动进程 for url in links: p = multiprocessing.Process(target=worker, args=(url, user_agents, result_queue)) processes.append(p) p.start() # 等待所有进程完成 for p in processes: p.join() # 从队列中提取结果 valid_records = [] while not result_queue.empty(): record = result_queue.get() if record: valid_records.append(record) print(f"成功解析 {len(valid_records)} 条商品数据")
注意事项
- 反爬规避:多进程并发请求易触发目标网站反爬,建议添加随机请求延迟(如
time.sleep(random.uniform(1,3)))、使用代理IP池,避免封禁IP。 - 数据序列化:
parse()方法返回的record必须是可序列化的(如字典、列表、基本类型),否则进程间无法传递数据。 - 进程数控制:不要设置过大的进程数,否则会导致网络资源耗尽或触发反爬,建议根据目标网站的承受能力调整。
内容的提问来源于stack exchange,提问作者Volatil3
相关产品推荐
相关产品推荐

