ListCollectService.java 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  1. package space.anyi.service;
  2. import space.anyi.client.GzCmcClient;
  3. import space.anyi.cleaner.DetailCleaner;
  4. import space.anyi.config.Config;
  5. import space.anyi.dto.ChannelAllContentsResponse;
  6. import space.anyi.dto.NewsItem;
  7. import space.anyi.dto.NewsItemData;
  8. import space.anyi.entity.News;
  9. import space.anyi.mapper.NewsMapper;
  10. import java.time.LocalDateTime;
  11. import java.util.ArrayList;
  12. import java.util.HashSet;
  13. import java.util.List;
  14. import java.util.Set;
  15. /**
  16. * 阶段一:列表采集服务。
  17. *
  18. * <p>对每个配置关键词全量分页采集列表,按发布时间窗口过滤,通过 id 去重后
  19. * 仅插入缺失记录(content 置空,作为阶段二待补队列)。可重复运行,幂等。</p>
  20. */
  21. public class ListCollectService {
  22. private final GzCmcClient client;
  23. private final NewsMapper mapper;
  24. private final Config config;
  25. /**
  26. * @param client 接口客户端
  27. * @param mapper NewsMapper
  28. * @param config 运行时配置
  29. */
  30. public ListCollectService(GzCmcClient client, NewsMapper mapper, Config config) {
  31. this.client = client;
  32. this.mapper = mapper;
  33. this.config = config;
  34. }
  35. /**
  36. * 执行列表采集主流程。
  37. *
  38. * <p>外层遍历配置的全部 channel、内层遍历全部关键词(双层循环,组合数为
  39. * O(频道数 × 关键词数) = O(n²)),每个频道×关键词组合独立分页采集。某页重试后
  40. * 仍失败(通常是服务端 offset≥10000 的检索上限)时优雅停止该组合的采集,
  41. * 不中断整体运行。完成后打印各组合与合计新增数。</p>
  42. *
  43. * @throws Exception 网络或序列化等未预期异常
  44. */
  45. public void run() throws Exception {
  46. int totalInserted = 0;
  47. for (String channelId : config.channelIds()) {
  48. int channelInserted = 0;
  49. for (String keyword : config.keywords()) {
  50. int pageNum = 1;
  51. int keywordInserted = 0;
  52. while (true) {
  53. ChannelAllContentsResponse resp;
  54. try {
  55. resp = client.search(keyword, channelId, pageNum, config.pageSize());
  56. } catch (Exception e) {
  57. System.err.printf("[列表] 频道=%s 关键词=%s pageNum=%d 请求失败(重试后),疑似达到服务端检索上限,停止本组合采集: %s%n",
  58. channelId, keyword, pageNum, e.getMessage());
  59. break;
  60. }
  61. List<NewsItem> items = resp.getList();
  62. if (items == null || items.isEmpty()) {
  63. break;
  64. }
  65. List<News> batch = new ArrayList<>();
  66. for (NewsItem item : items) {
  67. News news = toNews(item);
  68. if (news != null) {
  69. batch.add(news);
  70. }
  71. }
  72. keywordInserted += insertMissing(batch);
  73. int pages = resp.getPages() == null ? 1 : resp.getPages();
  74. System.out.printf("[列表] 频道=%s 关键词=%s pageNum=%d/%d 命中窗口=%d%n",
  75. channelId, keyword, pageNum, pages, batch.size());
  76. if (pageNum >= pages) {
  77. break;
  78. }
  79. pageNum++;
  80. }
  81. channelInserted += keywordInserted;
  82. System.out.printf("[列表] 频道=%s 关键词=%s 新增入库=%d%n", channelId, keyword, keywordInserted);
  83. }
  84. totalInserted += channelInserted;
  85. System.out.printf("[列表] 频道=%s 新增入库=%d%n", channelId, channelInserted);
  86. }
  87. System.out.printf("[列表] 全部完成, 新增入库合计=%d%n", totalInserted);
  88. }
  89. /**
  90. * 将列表元素映射为实体(含窗口过滤)。
  91. *
  92. * @param item 列表元素
  93. * @return News 实体;数据缺失、时间解析失败或不在发布窗口内时返回 null
  94. */
  95. private News toNews(NewsItem item) {
  96. NewsItemData data = item.getData();
  97. if (data == null || data.getId() == null) {
  98. return null;
  99. }
  100. LocalDateTime publishTime;
  101. try {
  102. publishTime = LocalDateTime.parse(data.getPublishTime(), Config.TIME_FORMATTER);
  103. } catch (Exception e) {
  104. return null;
  105. }
  106. if (publishTime.isBefore(config.publishStart()) || publishTime.isAfter(config.publishEnd())) {
  107. return null;
  108. }
  109. News news = new News();
  110. news.setId(data.getId());
  111. news.setTitle(DetailCleaner.titleClean(data.getTitle()));
  112. news.setUrl(data.getUrl());
  113. news.setPublishTime(publishTime);
  114. news.setChannelId(data.getChannelId());
  115. news.setChannelName(data.getChannelName());
  116. news.setEditor(data.getUserName());
  117. return news;
  118. }
  119. /**
  120. * 批量去重并插入缺失记录。
  121. *
  122. * <p>先用本批 id 查询库中已存在 id,仅插入不存在的记录,保证可重复运行。</p>
  123. *
  124. * @param batch 本页映射出的实体(已通过窗口过滤)
  125. * @return 实际新增行数
  126. */
  127. private int insertMissing(List<News> batch) {
  128. if (batch.isEmpty()) {
  129. return 0;
  130. }
  131. Set<String> ids = new HashSet<>();
  132. for (News n : batch) {
  133. ids.add(n.getId());
  134. }
  135. List<News> existing = mapper.selectBatchIds(ids);
  136. Set<String> existingIds = new HashSet<>();
  137. for (News n : existing) {
  138. existingIds.add(n.getId());
  139. }
  140. int inserted = 0;
  141. for (News n : batch) {
  142. if (existingIds.add(n.getId())) {
  143. mapper.insert(n);
  144. inserted++;
  145. }
  146. }
  147. return inserted;
  148. }
  149. }