package space.anyi.service; import space.anyi.client.GzCmcClient; import space.anyi.cleaner.DetailCleaner; import space.anyi.config.Config; import space.anyi.dto.ChannelAllContentsResponse; import space.anyi.dto.NewsItem; import space.anyi.dto.NewsItemData; import space.anyi.entity.News; import space.anyi.mapper.NewsMapper; import java.time.LocalDateTime; import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; /** * 阶段一:列表采集服务。 * *

对每个配置关键词全量分页采集列表,按发布时间窗口过滤,通过 id 去重后 * 仅插入缺失记录(content 置空,作为阶段二待补队列)。可重复运行,幂等。

*/ public class ListCollectService { private final GzCmcClient client; private final NewsMapper mapper; private final Config config; /** * @param client 接口客户端 * @param mapper NewsMapper * @param config 运行时配置 */ public ListCollectService(GzCmcClient client, NewsMapper mapper, Config config) { this.client = client; this.mapper = mapper; this.config = config; } /** * 执行列表采集主流程。 * *

外层遍历配置的全部 channel、内层遍历全部关键词(双层循环,组合数为 * O(频道数 × 关键词数) = O(n²)),每个频道×关键词组合独立分页采集。某页重试后 * 仍失败(通常是服务端 offset≥10000 的检索上限)时优雅停止该组合的采集, * 不中断整体运行。完成后打印各组合与合计新增数。

* * @throws Exception 网络或序列化等未预期异常 */ public void run() throws Exception { int totalInserted = 0; for (String channelId : config.channelIds()) { int channelInserted = 0; for (String keyword : config.keywords()) { int pageNum = 1; int keywordInserted = 0; while (true) { ChannelAllContentsResponse resp; try { resp = client.search(keyword, channelId, pageNum, config.pageSize()); } catch (Exception e) { System.err.printf("[列表] 频道=%s 关键词=%s pageNum=%d 请求失败(重试后),疑似达到服务端检索上限,停止本组合采集: %s%n", channelId, keyword, pageNum, e.getMessage()); break; } List items = resp.getList(); if (items == null || items.isEmpty()) { break; } List batch = new ArrayList<>(); for (NewsItem item : items) { News news = toNews(item); if (news != null) { batch.add(news); } } keywordInserted += insertMissing(batch); int pages = resp.getPages() == null ? 1 : resp.getPages(); System.out.printf("[列表] 频道=%s 关键词=%s pageNum=%d/%d 命中窗口=%d%n", channelId, keyword, pageNum, pages, batch.size()); if (pageNum >= pages) { break; } pageNum++; } channelInserted += keywordInserted; System.out.printf("[列表] 频道=%s 关键词=%s 新增入库=%d%n", channelId, keyword, keywordInserted); } totalInserted += channelInserted; System.out.printf("[列表] 频道=%s 新增入库=%d%n", channelId, channelInserted); } System.out.printf("[列表] 全部完成, 新增入库合计=%d%n", totalInserted); } /** * 将列表元素映射为实体(含窗口过滤)。 * * @param item 列表元素 * @return News 实体;数据缺失、时间解析失败或不在发布窗口内时返回 null */ private News toNews(NewsItem item) { NewsItemData data = item.getData(); if (data == null || data.getId() == null) { return null; } LocalDateTime publishTime; try { publishTime = LocalDateTime.parse(data.getPublishTime(), Config.TIME_FORMATTER); } catch (Exception e) { return null; } if (publishTime.isBefore(config.publishStart()) || publishTime.isAfter(config.publishEnd())) { return null; } News news = new News(); news.setId(data.getId()); news.setTitle(DetailCleaner.titleClean(data.getTitle())); news.setUrl(data.getUrl()); news.setPublishTime(publishTime); news.setChannelId(data.getChannelId()); news.setChannelName(data.getChannelName()); news.setEditor(data.getUserName()); return news; } /** * 批量去重并插入缺失记录。 * *

先用本批 id 查询库中已存在 id,仅插入不存在的记录,保证可重复运行。

* * @param batch 本页映射出的实体(已通过窗口过滤) * @return 实际新增行数 */ private int insertMissing(List batch) { if (batch.isEmpty()) { return 0; } Set ids = new HashSet<>(); for (News n : batch) { ids.add(n.getId()); } List existing = mapper.selectBatchIds(ids); Set existingIds = new HashSet<>(); for (News n : existing) { existingIds.add(n.getId()); } int inserted = 0; for (News n : batch) { if (existingIds.add(n.getId())) { mapper.insert(n); inserted++; } } return inserted; } }