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

如何将Scrapy响应输出保存到集群上的Minio Pod

解决Scrapy爬取内容存储到Minio的问题

核心方案:利用Scrapy FEEDS的S3兼容特性

Minio完全兼容S3协议,Scrapy的FEEDS配置原生支持S3存储,无需依赖全局变量这类绕路方案,直接配置即可实现爬取的CSV文件自动上传到Minio。

步骤1:安装依赖

Scrapy的S3 FEED依赖boto3和botocore,先执行安装:

pip install boto3 botocore

步骤2:配置Scrapy自定义设置

在爬虫类中修改custom_settings,添加S3兼容配置及FEEDS存储路径:

custom_settings = {
    # FEEDS配置,路径采用s3://<桶名>/<文件名>格式
    'FEEDS': {
        's3://your-minio-bucket/output.csv': {
            'format': 'csv',
            'overwrite': True,
            'fields': ['Text Paragraph']  # 可选:指定导出字段,避免冗余数据
        }
    },
    # Minio的S3兼容配置
    'AWS_ACCESS_KEY_ID': 'your-minio-access-key',
    'AWS_SECRET_ACCESS_KEY': 'your-minio-secret-key',
    # K8s集群内用Minio服务名+端口,例如http://minio-service:9000
    'AWS_ENDPOINT_URL': 'http://minio-pod-address:9000',
    'AWS_S3_SIGNATURE_VERSION': 's3v4',  # Minio默认支持的签名版本
    'AWS_REGION_NAME': ''  # 禁用AWS区域检查,Minio无需配置区域
}

步骤3:修正爬虫parse函数

移除无用的全局变量,直接返回清理后的段落数据,Scrapy会自动将每个yield的条目作为CSV的一行:

def parse(self, response):
    paragraphs = response.css('p')
    for paragraph in paragraphs:
        # 清理HTML标签与特殊字符
        clean_text = paragraph.get().replace('<p>', '').replace('</p>', '').replace('\xa0', '').strip()
        yield {
            'Text Paragraph': clean_text
        }

步骤4:简化handle函数

FEEDS会自动完成文件存储,无需手动处理全局变量和Minio上传操作,修改后的handle函数:

def handle(inputs: Dict, minio: Minio, log, **kwargs):
    with minio.get_object(input_bucket, "form.json") as response:
        form = json.loads(response.data.decode('utf-8'))
        url = form.get("select-url")
        spiderToCrawl = getSpiderToCrawl(url)
        logger.info(f"This is the url we are crawling: {url}")
        
        process = CrawlerProcess()
        process.crawl(spiderToCrawl, start_urls=[url])
        process.start()

备选方案:自定义Minio Pipeline(适用于更灵活的场景)

如果需要实时处理每条爬取数据,可编写自定义Pipeline:

from scrapy.exceptions import DropItem
import io
import csv
from minio import Minio

class MinioPipeline:
    def __init__(self, minio_endpoint, minio_access_key, minio_secret_key, minio_bucket):
        self.minio_endpoint = minio_endpoint
        self.minio_access_key = minio_access_key
        self.minio_secret_key = minio_secret_key
        self.minio_bucket = minio_bucket
        self.items = []
        self.fieldnames = ['Text Paragraph']

    @classmethod
    def from_crawler(cls, crawler):
        return cls(
            minio_endpoint=crawler.settings.get('MINIO_ENDPOINT'),
            minio_access_key=crawler.settings.get('MINIO_ACCESS_KEY'),
            minio_secret_key=crawler.settings.get('MINIO_SECRET_KEY'),
            minio_bucket=crawler.settings.get('MINIO_BUCKET')
        )

    def open_spider(self, spider):
        # 初始化Minio客户端
        self.client = Minio(
            self.minio_endpoint,
            access_key=self.minio_access_key,
            secret_key=self.minio_secret_key,
            secure=False  # 根据Minio配置选择True/False
        )
        # 检查桶是否存在,不存在则创建
        if not self.client.bucket_exists(self.minio_bucket):
            self.client.make_bucket(self.minio_bucket)

    def process_item(self, item, spider):
        self.items.append(item)
        return item

    def close_spider(self, spider):
        # 将数据写入CSV并上传到Minio
        output = io.StringIO()
        writer = csv.DictWriter(output, fieldnames=self.fieldnames)
        writer.writeheader()
        writer.writerows(self.items)
        
        csv_data = output.getvalue().encode('utf-8')
        self.client.put_object(
            self.minio_bucket,
            'output.csv',
            io.BytesIO(csv_data),
            length=len(csv_data),
            content_type='text/csv'
        )

随后在Scrapy全局配置中启用该Pipeline:

ITEM_PIPELINES = {
    'your_project.pipelines.MinioPipeline': 300,
}

# Minio全局配置
MINIO_ENDPOINT = 'minio-service:9000'
MINIO_ACCESS_KEY = 'your-access-key'
MINIO_SECRET_KEY = 'your-secret-key'
MINIO_BUCKET = 'your-bucket-name'

关键注意事项

  • K8s集群内访问Minio时,确保爬虫Pod能解析Minio的服务名(使用ClusterIP服务)或直接访问Pod IP。
  • 确认Minio密钥拥有目标桶的读写权限。
  • FEEDS方式适合批量生成文件后上传,Pipeline方式适合实时处理大数据量场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 02:22:45