You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

关键改进点

  1. Spider的数据关联:在parse方法里用next()精准匹配当前请求url对应的JSON条目,确保每个爬取页面只生成一个正确关联的Item,彻底解决重复循环的问题。
  2. SQL查询修正:使用参数化查询WHERE url = ?,既避免SQL注入风险,又能精准匹配当前Item的url。
  3. 空值安全处理:检查查询结果是否为None,避免索引错误;爬取页面数据时增加空值判断,用JSON里的默认值作为 fallback。
  4. 日志辅助调试:加入日志记录,方便跟踪数据插入/更新的执行情况。
  5. 规范表结构:补充完整建表语句,将url设为主键,从根源避免重复插入。

内容的提问来源于stack exchange,提问作者anfield

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 08:52:53