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;
}
}