在无Pipeline模块的Dagster中连接Web Scraper与Asset
在Dagster中整合Scrapy爬虫与Asset
你遇到的核心问题是不能直接调用Scrapy爬虫的parse方法——因为parse依赖Scrapy框架生成的response对象,必须通过Scrapy的执行器来完整运行爬虫流程。下面是修正后的可运行代码,以及关键逻辑说明:
完整实现代码
import scrapy from scrapy.crawler import CrawlerProcess from dagster import asset, AssetExecutionContext class MySpider(scrapy.Spider): name = 'headless' def start_requests(self): urls = ['http://google.com'] # 替换为目标URL for url in urls: yield scrapy.Request(url=url, callback=self.parse) def parse(self, response): headlines = response.css('h1::text').getall() yield {'headlines': headlines} @asset() def scraped_headlines(context: AssetExecutionContext): # 存储爬虫输出结果的容器 crawl_results = [] # 自定义Pipeline:捕获爬虫生成的Item并存入容器 class ResultCollectorPipeline: def process_item(self, item, spider): crawl_results.append(item) return item # 初始化Scrapy爬虫进程,配置Pipeline和日志级别 process = CrawlerProcess(settings={ 'ITEM_PIPELINES': { '__main__.ResultCollectorPipeline': 100, # 优先级数字越小越先执行 }, 'LOG_LEVEL': 'ERROR', # 禁用冗余日志,避免干扰Dagster日志 }) # 启动爬虫并等待执行完成 process.crawl(MySpider) process.start() # 提取并返回标题列表(适配单URL场景,多URL需遍历crawl_results) return crawl_results[0]['headlines'] if crawl_results else []
关键逻辑说明
- 爬虫执行方式:必须使用
CrawlerProcess启动爬虫,它会处理请求发送、响应接收、parse方法调用等完整流程,直接调用spider.parse()无法生成有效的response对象。 - 结果捕获:通过自定义的
ResultCollectorPipeline拦截爬虫输出的Item,将其存入内存列表crawl_results,这样在爬虫执行完成后就能获取到解析后的标题数据。 - Asset整合:在Dagster的
@asset装饰函数中完成爬虫的启动、结果收集和返回,最终scraped_headlines作为可被Dagster追踪的资产,可用于后续的下游任务。
额外调整建议
- 若爬取多个URL,
crawl_results会包含多个Item对象,可根据需求将所有标题合并或按URL分组返回。 - 可在
CrawlerProcess的settings中添加Scrapy的其他配置(如DOWNLOAD_DELAY、USER_AGENT等),优化爬虫行为。 - 若爬虫需要复杂前置操作(如登录、Cookie处理),直接在
MySpider类中实现即可,不影响与Dagster Asset的整合。
内容的提问来源于stack exchange,提问作者marcel
相关产品推荐
相关产品推荐

