feat(ZmswpmService): enhance API data fetching with retry logic and update data synchronization methods

This commit is contained in:
2025-10-22 21:55:24 +08:00
parent 527560aec3
commit 35383cd3a3

View File

@@ -56,6 +56,8 @@ public class ZmswpmService extends BaseService<ZmswpmEntity, ZmswpmRepository> {
private static final String API_URL =
"https://www.chinacoalst.cn/electronicbusiness/shoppingApi/bidttargpage";
private static final int PAGE_SIZE = 50;
private static final int FETCH_API_MAX_RETRY = 3;
private static final long FETCH_API_RETRY_DELAY_MS = 3000L;
private static final DateTimeFormatter DATE_TIME_FORMATTER =
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm");
@Autowired private ZmswpmRepository repository;
@@ -216,11 +218,14 @@ public class ZmswpmService extends BaseService<ZmswpmEntity, ZmswpmRepository> {
long totalCount = repository.count();
if (totalCount == 0) {
// 数据库为空,调用接口直到没有数据返回
syncAllData();
// 数据库为空,采集最近一个月的数据
LocalDate now = LocalDate.now();
LocalDate oneMonthAgo = now.minusMonths(1);
log.info("数据库为空,采集{}到{}之间的数据", oneMonthAgo, now);
syncRecentData(oneMonthAgo, now);
} else {
// 数据库不为空,同步当前天和前天的数据
syncRecentData(LocalDate.now(), LocalDate.now().minusDays(7));
// 数据库不为空,同步当前天和前天的数据
syncRecentData(LocalDate.now().minusDays(7), LocalDate.now());
}
log.info("煤炭商务网拍卖数据同步完成");
@@ -230,7 +235,12 @@ public class ZmswpmService extends BaseService<ZmswpmEntity, ZmswpmRepository> {
}
}
/** 同步所有数据(数据库为空时使用) */
/**
* 同步所有数据(数据库为空时使用)
*
* @deprecated 已由 syncRecentData 替代,在 syncAuctionData 中采集最近一个月的数据
*/
@Deprecated(since = "2.0", forRemoval = true)
private void syncAllData() {
int currentPage = 1;
boolean hasMoreData = true;
@@ -311,148 +321,165 @@ public class ZmswpmService extends BaseService<ZmswpmEntity, ZmswpmRepository> {
/** 从API获取数据支持按标的名称和时间范围查询 */
private List<Map<String, Object>> fetchDataFromApi(
int page, String bidname, String begintime, String endtime) {
for (int attempt = 1; attempt <= FETCH_API_MAX_RETRY; attempt++) {
ProxyIp currentProxy = null;
RestTemplate currentRestTemplate = restTemplate;
ProxyIp currentProxy = null;
RestTemplate currentRestTemplate = restTemplate;
try {
// 获取可用的代理IP
Optional<ProxyIp> proxyOptional = proxyIpService.getAvailableProxyIp();
if (proxyOptional.isPresent()) {
currentProxy = proxyOptional.get();
currentRestTemplate = createProxyRestTemplate(currentProxy);
log.info("使用代理IP进行API请求{}:{}", currentProxy.getIp(), currentProxy.getPort());
} else {
log.warn("未获取到可用代理IP使用直接连接");
}
// 设置请求头
HttpHeaders headers = new HttpHeaders();
headers.set("Accept", "*/*");
headers.set("Accept-Language", "zh-CN,zh;q=0.9");
headers.set("Cache-Control", "no-cache");
headers.set("Connection", "keep-alive");
headers.set("host", "www.chinacoalst.cn");
headers.set("x-auth-token", "");
headers.set("sec-ch-ua-platform", "\"Windows\"");
headers.set(
"sec-ch-ua",
"\"Chromium\";v=\"140\", \"Not=A?Brand\";v=\"24\", \"Google Chrome\";v=\"140\"");
headers.set("sec-ch-ua-mobile", "?0");
headers.set("language-type-code", "null");
headers.set("id", "9db233983befbae5");
headers.set(
"user-agent",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36");
headers.set("token", "null");
headers.set("sec-fetch-site", "same-origin");
headers.set("sec-fetch-mode", "cors");
headers.set("sec-fetch-dest", "empty");
headers.set("referer", "https://www.chinacoalst.cn/");
headers.set("accept-encoding", "gzip, deflate, br, zstd");
headers.set("priority", "u=1, i");
headers.set(
"Cookie",
"acw_tc=0bd17c0e17606243767095896e3ece80da4344e5f095bfd2618afec87f6899");
// 构建URL参数
StringBuilder urlBuilder = new StringBuilder(API_URL);
urlBuilder.append("?page=").append(page);
urlBuilder.append("&limit=").append(PAGE_SIZE);
urlBuilder.append("&productid=");
urlBuilder.append("&delivery=");
urlBuilder.append("&outbidstate=");
urlBuilder.append("&coalType=");
urlBuilder.append("&splitstate=");
urlBuilder.append("&tenantCode=");
urlBuilder.append("&type=1");
if (begintime != null && !begintime.isEmpty()) {
urlBuilder.append("&begintime=").append(begintime);
} else {
urlBuilder.append("&begintime=");
}
if (endtime != null && !endtime.isEmpty()) {
urlBuilder.append("&endtime=").append(endtime);
} else {
urlBuilder.append("&endtime=");
}
if (bidname != null && !bidname.trim().isEmpty()) {
urlBuilder.append("&bidname=").append(bidname.trim());
} else {
urlBuilder.append("&bidname=");
}
urlBuilder.append("&mineName=");
HttpEntity<String> request = new HttpEntity<>(headers);
String url = urlBuilder.toString();
// 发送请求
ResponseEntity<String> response =
currentRestTemplate.getForEntity(url, String.class, request);
if (response.getStatusCode() == HttpStatus.OK) {
// 标记代理IP使用成功
if (currentProxy != null) {
proxyIpService.markProxyAsUsed(currentProxy.getIp(), currentProxy.getPort());
log.debug("代理IP {}:{} 请求成功", currentProxy.getIp(), currentProxy.getPort());
try {
// 获取可用的代理IP
Optional<ProxyIp> proxyOptional = proxyIpService.getAvailableProxyIp();
if (proxyOptional.isPresent()) {
currentProxy = proxyOptional.get();
currentRestTemplate = createProxyRestTemplate(currentProxy);
log.info(
"使用代理IP进行API请求{}:{},第{}次尝试",
currentProxy.getIp(),
currentProxy.getPort(),
attempt);
} else {
log.warn("未获取到可用代理IP使用直接连接第{}次尝试", attempt);
}
String responseBody = response.getBody();
Map<String, Object> responseMap =
objectMapper.readValue(
responseBody, new TypeReference<Map<String, Object>>() {});
// 设置请求头
HttpHeaders headers = new HttpHeaders();
headers.set("Accept", "*/*");
headers.set("Accept-Language", "zh-CN,zh;q=0.9");
headers.set("Cache-Control", "no-cache");
headers.set("Connection", "keep-alive");
headers.set("host", "www.chinacoalst.cn");
headers.set("x-auth-token", "");
headers.set("sec-ch-ua-platform", "\"Windows\"");
headers.set(
"sec-ch-ua",
"\"Chromium\";v=\"140\", \"Not=A?Brand\";v=\"24\", \"Google Chrome\";v=\"140\"");
headers.set("sec-ch-ua-mobile", "?0");
headers.set("language-type-code", "null");
headers.set("id", "9db233983befbae5");
headers.set(
"user-agent",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36");
headers.set("token", "null");
headers.set("sec-fetch-site", "same-origin");
headers.set("sec-fetch-mode", "cors");
headers.set("sec-fetch-dest", "empty");
headers.set("referer", "https://www.chinacoalst.cn/");
headers.set("accept-encoding", "gzip, deflate, br, zstd");
headers.set("priority", "u=1, i");
headers.set(
"Cookie",
"acw_tc=0bd17c0e17606243767095896e3ece80da4344e5f095bfd2618afec87f6899");
Integer code = (Integer) responseMap.get("code");
if (code != null && code == 0) {
Object dataObj = responseMap.get("data");
if (dataObj instanceof Map) {
@SuppressWarnings("unchecked")
Map<String, Object> dataMap = (Map<String, Object>) dataObj;
Object listObj = dataMap.get("list");
if (listObj instanceof List) {
// 构建URL参数
StringBuilder urlBuilder = new StringBuilder(API_URL);
urlBuilder.append("?page=").append(page);
urlBuilder.append("&limit=").append(PAGE_SIZE);
urlBuilder.append("&productid=");
urlBuilder.append("&delivery=");
urlBuilder.append("&outbidstate=");
urlBuilder.append("&coalType=");
urlBuilder.append("&splitstate=");
urlBuilder.append("&tenantCode=");
urlBuilder.append("&type=1");
if (begintime != null && !begintime.isEmpty()) {
urlBuilder.append("&begintime=").append(begintime);
} else {
urlBuilder.append("&begintime=");
}
if (endtime != null && !endtime.isEmpty()) {
urlBuilder.append("&endtime=").append(endtime);
} else {
urlBuilder.append("&endtime=");
}
if (bidname != null && !bidname.trim().isEmpty()) {
urlBuilder.append("&bidname=").append(bidname.trim());
} else {
urlBuilder.append("&bidname=");
}
urlBuilder.append("&mineName=");
HttpEntity<String> requestEntity = new HttpEntity<>(headers);
String url = urlBuilder.toString();
// 发送请求
ResponseEntity<String> response =
currentRestTemplate.exchange(url, HttpMethod.GET, requestEntity, String.class);
if (response.getStatusCode() == HttpStatus.OK) {
// 标记代理IP使用成功
if (currentProxy != null) {
proxyIpService.markProxyAsUsed(currentProxy.getIp(), currentProxy.getPort());
log.debug("代理IP {}:{} 请求成功", currentProxy.getIp(), currentProxy.getPort());
}
String responseBody = response.getBody();
Map<String, Object> responseMap =
objectMapper.readValue(
responseBody, new TypeReference<Map<String, Object>>() {});
Integer code = (Integer) responseMap.get("code");
if (code != null && code == 0) {
Object dataObj = responseMap.get("data");
if (dataObj instanceof Map) {
@SuppressWarnings("unchecked")
List<Map<String, Object>> dataList =
(List<Map<String, Object>>) listObj;
return dataList;
Map<String, Object> dataMap = (Map<String, Object>) dataObj;
Object listObj = dataMap.get("list");
if (listObj instanceof List) {
@SuppressWarnings("unchecked")
List<Map<String, Object>> dataList =
(List<Map<String, Object>>) listObj;
return dataList;
}
}
}
log.warn("API返回成功状态但数据结构不符合预期第{}次尝试", attempt);
} else {
log.warn("API请求返回非200状态码{}", response.getStatusCode());
// 如果是代理相关的错误标记代理IP为不可用
if (currentProxy != null
&& (response.getStatusCode().is4xxClientError()
|| response.getStatusCode().is5xxServerError())) {
proxyIpService.markProxyAsUnavailable(
currentProxy.getIp(), currentProxy.getPort());
log.warn(
"代理IP {}:{} 可能不可用,状态码:{}",
currentProxy.getIp(),
currentProxy.getPort(),
response.getStatusCode());
}
}
} else {
log.warn("API请求返回非200状态码{}", response.getStatusCode());
// 如果是代理相关的错误标记代理IP为不可用
if (currentProxy != null
&& (response.getStatusCode().is4xxClientError()
|| response.getStatusCode().is5xxServerError())) {
} catch (Exception e) {
log.error("API获取数据失败,页码:{},第{}次尝试", page, attempt, e);
// 如果使用了代理IP且请求失败标记为不可用
if (currentProxy != null) {
proxyIpService.markProxyAsUnavailable(
currentProxy.getIp(), currentProxy.getPort());
log.warn(
"代理IP {}:{} 可能不可用,状态码{}",
"代理IP {}:{} 请求失败,标记为不可用{}",
currentProxy.getIp(),
currentProxy.getPort(),
response.getStatusCode());
e.getMessage());
}
}
return null;
} catch (Exception e) {
log.error("从API获取数据失败页码{}", page, e);
// 如果使用了代理IP且请求失败标记为不可用
if (currentProxy != null) {
proxyIpService.markProxyAsUnavailable(currentProxy.getIp(), currentProxy.getPort());
log.warn(
"代理IP {}:{} 请求失败,标记为不可用:{}",
currentProxy.getIp(),
currentProxy.getPort(),
e.getMessage());
if (attempt < FETCH_API_MAX_RETRY) {
log.info("将在{}毫秒后重试({}/{})", FETCH_API_RETRY_DELAY_MS, attempt, FETCH_API_MAX_RETRY);
try {
Thread.sleep(FETCH_API_RETRY_DELAY_MS);
} catch (InterruptedException interruptedException) {
Thread.currentThread().interrupt();
log.warn("重试等待被中断", interruptedException);
break;
}
}
return null;
}
return null;
}
/** 保存数据到数据库 */