如何将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
相关产品推荐
相关产品推荐

