Scrapy中如何从JSON传递额外数据至Pipeline实现DB增改
解决Scrapy中从JSON传递数据到Pipeline并实现增改逻辑的问题
我来帮你修正这两个文件里的问题,主要是Spider里的数据关联逻辑错误,还有Pipeline里的SQL查询和空值处理问题。
原代码的核心问题
- Spider的循环逻辑错误:
parse方法里嵌套了两个for循环,每个请求会遍历所有start_urls和所有JSON里的itemdata,导致生成大量重复的Item,而且无法正确将当前爬取的页面和JSON里对应的条目关联起来。 - Pipeline的SQL查询错误:
get_data里的SQL语句select url, new_price from price_monitor WHERE url=url是无效的,这里的url=url永远为真,会返回所有数据,应该用参数化查询匹配当前Item的url。 - 空值处理缺失:如果数据库里没有对应url的记录,
rows = self.curr.fetchone()会返回None,直接访问rows[0]会抛出IndexError。
修正后的myspider.py
import scrapy import json from ..items import AmazonItem class MySpider(scrapy.Spider): name = 'price_monitor' itemdatalist = [] def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) # 读取JSON数据,建议后续改用相对路径或配置文件,避免硬编码绝对路径 with open('C:\\Users\\Documents\\python_virtual\\price_monitor\\price_monitor\\products.json') as f: data = json.load(f) self.itemdatalist = data['itemdata'] # 生成start_urls self.start_urls = [item['url'] for item in self.itemdatalist] def parse(self, response): # 找到当前请求url对应的JSON条目 current_item_data = next((item for item in self.itemdatalist if item['url'] == response.url), None) if not current_item_data: self.logger.warning(f"No matching data found for URL: {response.url}") return scrapeitem = AmazonItem() # 爬取页面数据,增加空值判断 title = response.css('span#productTitle::text').extract_first() scrapeitem['title'] = title.strip() if title else current_item_data.get('title', '') price = response.css('span#priceblock_ourprice::text').extract_first() scrapeitem['price'] = price if price else '' # 从JSON里获取额外关联数据 scrapeitem['url'] = current_item_data['url'] scrapeitem['name'] = current_item_data['name'] scrapeitem['email'] = current_item_data['email'] yield scrapeitem
修正后的pipelines.py
import sqlite3 class PriceMonitorPipeline(object): def __init__(self): self.create_connection() self.create_table() def create_connection(self): self.conn = sqlite3.connect("price_monitor.db") self.curr = self.conn.cursor() def create_table(self): # 补充完整建表语句,确保url为主键避免重复 self.curr.execute(""" CREATE TABLE IF NOT EXISTS price_monitor ( url TEXT PRIMARY KEY, title TEXT, name TEXT, email TEXT, old_price TEXT, new_price TEXT ) """) self.conn.commit() def process_item(self, item, spider): self.get_data(item) return item def get_data(self, item): """检查当前url对应的记录是否存在""" # 正确的参数化查询,避免SQL注入并精准匹配 self.curr.execute("""SELECT url, new_price FROM price_monitor WHERE url = ?""", (item['url'],)) rows = self.curr.fetchone() if rows: # 存在记录则执行更新 rows_url, new_price = rows self.set_data_update(item, rows_url, new_price) spider.logger.info(f"Updated record for URL: {item['url']}") else: # 不存在记录则执行插入 self.set_insert_data(item) spider.logger.info(f"Inserted new record for URL: {item['url']}") def set_insert_data(self, item): self.curr.execute(""" INSERT INTO price_monitor (url, title, name, email, old_price, new_price) VALUES (?, ?, ?, ?, ?, ?) """, ( item['url'], item['title'], item['name'], item['email'], item['price'], item['price'] )) self.conn.commit() def set_data_update(self, item, rows_url, new_price): # 将原new_price转为old_price,更新为新爬取的price self.curr.execute(""" UPDATE price_monitor SET old_price = ?, new_price = ? WHERE url = ? """, (new_price, item['price'], rows_url)) self.conn.commit()
关键改进点
- Spider的数据关联:在
parse方法里用next()精准匹配当前请求url对应的JSON条目,确保每个爬取页面只生成一个正确关联的Item,彻底解决重复循环的问题。 - SQL查询修正:使用参数化查询
WHERE url = ?,既避免SQL注入风险,又能精准匹配当前Item的url。 - 空值安全处理:检查查询结果是否为
None,避免索引错误;爬取页面数据时增加空值判断,用JSON里的默认值作为 fallback。 - 日志辅助调试:加入日志记录,方便跟踪数据插入/更新的执行情况。
- 规范表结构:补充完整建表语句,将
url设为主键,从根源避免重复插入。
内容的提问来源于stack exchange,提问作者anfield
相关产品推荐
相关产品推荐

