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

StormCrawler多域名爬取完成时触发动作的技术咨询

在StormCrawler多域名爬取场景中判断单个域名爬取完成的实现方案

在StormCrawler+Elasticsearch的架构下,并没有内置的“域名爬取完成”判断组件,因为爬取是包含周期性重访的持续过程。我们需要结合ES的URL状态存储,自定义逻辑来实现这个需求,核心思路是跟踪每个域名的待爬URL数量,当数量归零时触发动作。下面是具体的实现方式:

一、核心判断逻辑

StormCrawler会将所有URL的状态(DISCOVERED/FETCH_SCHEDULED/FETCHED/ERROR等)存储在Elasticsearch的URL索引中。我们可以通过以下规则判断单个域名的爬取完成:

  1. 针对目标域名,统计ES中状态为DISCOVERED(已发现待爬)或FETCH_SCHEDULED(已调度待爬)的URL数量
  2. 当该数量为0,且没有新的URL被发现时,即可认为该域名当前爬取周期完成(首次爬取或某次重访周期)

二、具体实现步骤

1. 自定义DomainCompletionBolt

实现一个自定义Bolt,订阅StatusUpdaterBolt的输出流,监听每个URL的处理完成事件,然后触发域名完成判断:

public class DomainCompletionBolt extends BaseRichBolt {
    private EsClient esClient;
    private OutputCollector collector;
    // 分布式场景建议用ES存储完成状态,避免本地缓存不一致
    private Map<String, Boolean> completedDomainCache = new ConcurrentHashMap<>();

    @Override
    public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
        this.esClient = EsClientFactory.getClient(topoConf);
    }

    @Override
    public void execute(Tuple input) {
        String url = input.getStringByField("url");
        String domain = URLUtil.getDomainName(url);
        String status = input.getStringByField("status");

        // 只处理已完成爬取的URL(成功或失败)
        if (Status.FETCHED.name().equals(status) || Status.ERROR.name().equals(status)) {
            long pendingUrlCount = countPendingUrlsForDomain(domain);
            // 待爬数为0且未触发过完成信号时,发送Tuple到目标Bolt
            if (pendingUrlCount == 0 && !completedDomainCache.getOrDefault(domain, false)) {
                collector.emit(new Values(domain, new Date()));
                completedDomainCache.put(domain, true);
                
                // 重访场景:可根据重访间隔设置定时器,到期后重置缓存标记
                // 或在Scheduler生成重访URL时发送信号重置标记
            }
        }
        collector.ack(input);
    }

    // 查询ES中目标域名的待爬URL数量
    private long countPendingUrlsForDomain(String domain) throws IOException {
        SearchRequest request = new SearchRequest("url");
        BoolQueryBuilder query = QueryBuilders.boolQuery()
                .must(QueryBuilders.termQuery("domain", domain))
                .must(QueryBuilders.termsQuery("status", 
                      Status.DISCOVERED.name(), Status.FETCH_SCHEDULED.name()));
        request.source().query(query).size(0);
        
        SearchResponse response = esClient.search(request, RequestOptions.DEFAULT);
        return response.getHits().getTotalHits().value;
    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(new Fields("domain", "completion_time"));
    }
}

2. 集成到Storm拓扑

将自定义Bolt加入拓扑,连接到StatusUpdaterBolt的输出流,并对接你的目标Bolt:

TopologyBuilder builder = new TopologyBuilder();
// ... 其他组件(Spout、FetchBolt、StatusUpdaterBolt等)的定义 ...

// 添加域名完成判断Bolt
builder.setBolt("domain-completion", new DomainCompletionBolt())
        .shuffleGrouping("status-updater");

// 连接到你需要发送Tuple的目标Bolt
builder.setBolt("target-action-bolt", new YourTargetBolt())
        .shuffleGrouping("domain-completion");

3. 处理重访场景的优化

如果需要支持周期性重访,要避免重复触发或遗漏重访周期的完成信号:

  • 在URL元数据中添加crawl_cycle字段,标记当前爬取周期(比如首次爬取、第1次重访等)
  • 查询待爬URL数量时,同时过滤crawl_cycle字段,确保判断的是当前周期的待爬数
  • 当Scheduler生成重访URL时,发送信号到DomainCompletionBolt,重置对应域名的完成标记

4. 性能优化建议

  • 批量查询:避免每个URL都触发ES查询,可以每隔10-30秒批量查询所有域名的待爬数
  • ES Watcher监控:如果使用ES X-Pack,可以配置Watcher监控URL索引,当某个域名的待爬数变为0时,主动发送通知到Storm拓扑
  • 分布式状态存储:用ES代替本地缓存存储域名完成状态,避免多Bolt实例间的缓存不一致问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:20:39