提交货源采集更新

This commit is contained in:
super
2026-05-10 01:00:08 +08:00
parent a4e4745921
commit 5c990e651e
23 changed files with 308 additions and 37 deletions

View File

@@ -9,7 +9,7 @@ public class SimilarAsinProperties {
private String cozeBaseUrl = "https://api.coze.cn";
private String cozeWorkflowPath = "/v1/workflow/run";
private String cozeWorkflowHistoryPath = "/v1/workflows/{workflow_id}/run_histories/{execute_id}";
private String cozeWorkflowId = "7632683471312355338";
private String cozeWorkflowId = "7635328462404583478";
private String cozeToken = "";
private int cozeBatchSize = 50;
private int cozeConnectTimeoutMillis = 10000;

View File

@@ -556,7 +556,6 @@ public class AppearancePatentTaskService {
: row.getResultFilename();
}
@Scheduled(cron = "${aiimage.appearance-patent.stale-finalize-cron:0 */2 * * * *}")
public void finalizeStaleTasks() {
if (transactionManager != null) {
LocalDateTime threshold = LocalDateTime.now().minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes()));

View File

@@ -1209,7 +1209,6 @@ public class BrandTaskService {
}
@Transactional
@org.springframework.scheduling.annotation.Scheduled(cron = "${aiimage.brand-progress.stale-check-cron:0 */2 * * * *}")
public void failStaleRunningTasks() {
DistributedJobLockService.LockHandle lockHandle = distributedJobLockService.tryLock("brand:stale-check", STALE_CHECK_LOCK_TTL);
if (lockHandle == null) {

View File

@@ -4,6 +4,8 @@ import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.nanri.aiimage.common.service.DistributedJobLockService;
import com.nanri.aiimage.config.DeleteBrandProgressProperties;
import com.nanri.aiimage.modules.appearancepatent.service.AppearancePatentTaskService;
import com.nanri.aiimage.modules.brand.service.BrandTaskService;
import com.nanri.aiimage.modules.productrisk.service.ProductRiskTaskCacheService;
import com.nanri.aiimage.modules.queryasin.service.QueryAsinTaskCacheService;
import com.nanri.aiimage.modules.queryasin.service.QueryAsinTaskService;
@@ -12,6 +14,7 @@ import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskCacheServi
import com.nanri.aiimage.modules.patroldelete.service.PatrolDeleteTaskService;
import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskCacheService;
import com.nanri.aiimage.modules.pricetrack.service.PriceTrackTaskService;
import com.nanri.aiimage.modules.similarasin.service.SimilarAsinTaskService;
import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskCacheService;
import com.nanri.aiimage.modules.shopmatch.service.ShopMatchTaskService;
import com.nanri.aiimage.modules.task.mapper.FileTaskMapper;
@@ -63,6 +66,9 @@ public class DeleteBrandStaleTaskService {
private final PatrolDeleteTaskCacheService patrolDeleteTaskCacheService;
private final QueryAsinTaskService queryAsinTaskService;
private final QueryAsinTaskCacheService queryAsinTaskCacheService;
private final BrandTaskService brandTaskService;
private final AppearancePatentTaskService appearancePatentTaskService;
private final SimilarAsinTaskService similarAsinTaskService;
private final DeleteBrandProgressProperties deleteBrandProgressProperties;
private final DistributedJobLockService distributedJobLockService;
private final TaskDistributedLockService taskDistributedLockService;
@@ -88,6 +94,9 @@ public class DeleteBrandStaleTaskService {
ShopMatchStaleCheckStats shopMatchStats = failStaleShopMatchTasks();
ShopMatchStaleCheckStats patrolDeleteStats = failStalePatrolDeleteTasks();
ShopMatchStaleCheckStats queryAsinStats = failStaleQueryAsinTasks();
runModuleStaleCheck("brand", brandTaskService::failStaleRunningTasks);
runModuleStaleCheck("appearance-patent", appearancePatentTaskService::finalizeStaleTasks);
runModuleStaleCheck("similar-asin", similarAsinTaskService::finalizeStaleTasks);
log.info("[stale-check] product-risk summary scanned={} finalized={} failed={} skipped={} elapsedMs={} thread={}",
stats.scannedTaskCount,
stats.finalizedTaskCount,
@@ -126,6 +135,23 @@ public class DeleteBrandStaleTaskService {
}
}
private void runModuleStaleCheck(String moduleName, Runnable action) {
long startedAt = System.currentTimeMillis();
try {
action.run();
log.info("[stale-check] {} delegated stale check completed elapsedMs={} thread={}",
moduleName,
System.currentTimeMillis() - startedAt,
Thread.currentThread().getName());
} catch (Exception ex) {
log.warn("[stale-check] {} delegated stale check failed elapsedMs={} msg={}",
moduleName,
System.currentTimeMillis() - startedAt,
ex.getMessage(),
ex);
}
}
private void failStaleDeleteBrandTasks() {
long minutes = Math.max(1L, deleteBrandProgressProperties.getHeartbeatTimeoutMinutes());
long initialMinutes = Math.max(minutes, deleteBrandProgressProperties.getDeleteBrandInitialTimeoutMinutes());

View File

@@ -210,6 +210,9 @@ public class SimilarAsinCozeClient {
Map<String, Object> body = new LinkedHashMap<>();
body.put("workflow_id", properties.getCozeWorkflowId());
body.put("parameters", parameters);
if (apiKey != null && !apiKey.isBlank()) {
body.put("api_key", apiKey.trim());
}
body.put("is_async", Boolean.TRUE);
log.info("[similar-asin] coze request url={} body={}",
joinUrl(properties.getCozeBaseUrl(), properties.getCozeWorkflowPath()),
@@ -256,20 +259,12 @@ public class SimilarAsinCozeClient {
}
private Map<String, Object> buildParameters(List<SimilarAsinResultRowDto> rows, String prompt, String apiKey) {
List<String> groupKeys = rows.stream().map(row -> nonBlank(row.getGroupKey(), rowKey(row))).toList();
List<String> rowIds = rows.stream().map(row -> nonBlank(row.getId(), "")).toList();
List<String> asins = rows.stream().map(row -> nonBlank(row.getAsin(), "")).toList();
List<String> countries = rows.stream().map(row -> nonBlank(row.getCountry(), "")).toList();
List<String> prices = rows.stream().map(row -> nonBlank(row.getPrice(), "")).toList();
List<String> titles = rows.stream().map(row -> nonBlank(row.getTitle(), row.getAsin())).toList();
List<String> urls = rows.stream().map(row -> nonBlank(row.getUrl(), "")).toList();
List<List<String>> urlLists = rows.stream().map(SimilarAsinResultRowDto::getUrls).toList();
Map<String, Object> parameters = new LinkedHashMap<>();
parameters.put("title_list", titles);
parameters.put("url_list", urls);
parameters.put("url_lists", urlLists);
parameters.put("items", buildItemObjects(rows, groupKeys, rowIds, asins, countries, prices, titles, urls, urlLists));
parameters.put("items", buildItemObjects(rows, asins, urls, urlLists));
parameters.put("prompt", prompt == null ? "" : prompt);
if (apiKey != null && !apiKey.isBlank()) {
parameters.put("api_key", apiKey.trim());
@@ -280,6 +275,10 @@ public class SimilarAsinCozeClient {
@SuppressWarnings("unchecked")
private Map<String, Object> maskCozeRequestBody(Map<String, Object> body) {
Map<String, Object> masked = new LinkedHashMap<>(body);
Object topLevelApiKey = masked.get("api_key");
if (topLevelApiKey instanceof String apiKeyText && !apiKeyText.isBlank()) {
masked.put("api_key", maskSecret(apiKeyText));
}
Object parametersObj = masked.get("parameters");
if (parametersObj instanceof Map<?, ?> parameters) {
Map<String, Object> maskedParameters = new LinkedHashMap<>((Map<String, Object>) parameters);
@@ -304,25 +303,14 @@ public class SimilarAsinCozeClient {
}
private List<Map<String, Object>> buildItemObjects(List<SimilarAsinResultRowDto> rows,
List<String> groupKeys,
List<String> rowIds,
List<String> asins,
List<String> countries,
List<String> prices,
List<String> titles,
List<String> urls,
List<List<String>> urlLists) {
List<Map<String, Object>> items = new ArrayList<>(rows.size());
for (int i = 0; i < rows.size(); i++) {
Map<String, Object> item = new LinkedHashMap<>();
item.put("group_key", groupKeys.get(i));
item.put("row_id", rowIds.get(i));
item.put("asin", asins.get(i));
item.put("country", countries.get(i));
item.put("price", prices.get(i));
item.put("title", titles.get(i));
item.put("url", urls.get(i));
item.put("urls", urlLists.get(i));
item.put("target_urls", urlLists.get(i));
items.add(item);
}

View File

@@ -616,7 +616,6 @@ public class SimilarAsinTaskService {
: row.getResultFilename();
}
@Scheduled(cron = "${aiimage.similar-asin.stale-finalize-cron:0 */2 * * * *}")
public void finalizeStaleTasks() {
if (transactionManager != null) {
LocalDateTime threshold = LocalDateTime.now().minusMinutes(Math.max(1, properties.getStaleTimeoutMinutes()));
@@ -1429,6 +1428,10 @@ public class SimilarAsinTaskService {
SimilarAsinCozeClient.CozeSubmitResponse submit = submitCozeWorkflowThrottled(batchRows, prompt, apiKey);
if (submit.immediateData() != null && !submit.immediateData().isBlank()) {
List<SimilarAsinResultRowDto> cozeRows = cozeClient.mergeRowsFromDataText(batchRows, submit.immediateData());
String emptyResultMessage = emptyCozeResultMessage(cozeRows, batchRows.size());
if (!emptyResultMessage.isBlank()) {
throw new IllegalStateException(emptyResultMessage);
}
mergeCozeRowsIntoChunk(task, chunk.getScopeHash(), chunk.getChunkIndex(), cozeRows, allRowsByBaseId);
return false;
}
@@ -1570,15 +1573,22 @@ public class SimilarAsinTaskService {
if (batchRows.isEmpty() && failureMessage.isBlank()) {
failureMessage = "Coze batch payload missing";
}
List<SimilarAsinResultRowDto> cozeRows;
if (failureMessage.isBlank()) {
cozeRows = cozeClient.mergeRowsFromDataText(batchRows, poll.resolvedPayloadText());
failureMessage = emptyCozeResultMessage(cozeRows, batchRows.size());
} else {
cozeRows = List.of();
}
if (!failureMessage.isBlank() && splitRetryFailedCozeBatchState(state, context, batchRows, failureMessage)) {
return;
}
if (!failureMessage.isBlank() && retryFailedCozeBatchState(state, context, batchRows, failureMessage)) {
return;
}
List<SimilarAsinResultRowDto> cozeRows = failureMessage.isBlank()
? cozeClient.mergeRowsFromDataText(batchRows, poll.resolvedPayloadText())
: cozeClient.markRowsFailed(batchRows, failureMessage);
if (!failureMessage.isBlank()) {
cozeRows = cozeClient.markRowsFailed(batchRows, failureMessage);
}
FileTaskEntity task = fileTaskMapper.selectById(state.getTaskId());
if (task != null) {
Map<String, List<SimilarAsinParsedRowVo>> allRowsByBaseId = loadAllRowsByBaseId(task);
@@ -1819,6 +1829,7 @@ public class SimilarAsinTaskService {
String normalized = normalize(failureMessage).toLowerCase(Locale.ROOT);
return normalized.contains("timeout")
|| normalized.contains("timed out")
|| normalized.contains("empty result rows")
|| normalized.contains("out of limit")
|| normalized.contains("execution limit")
|| normalized.contains("720712008")
@@ -1834,6 +1845,7 @@ public class SimilarAsinTaskService {
|| normalized.contains("retry later")
|| normalized.contains("timeout")
|| normalized.contains("timed out")
|| normalized.contains("empty result rows")
|| normalized.contains("out of limit")
|| normalized.contains("execution limit")
|| normalized.contains("702093018")
@@ -1846,6 +1858,21 @@ public class SimilarAsinTaskService {
|| normalized.contains("\u8c03\u7528\u8d85\u65f6");
}
private String emptyCozeResultMessage(List<SimilarAsinResultRowDto> cozeRows, int expectedRows) {
if (cozeRows == null || cozeRows.isEmpty()) {
return expectedRows > 0 ? "Coze async workflow returned empty result rows" : "";
}
long unresolved = cozeRows.stream()
.filter(row -> row != null
&& !hasResolvedCozeFields(row)
&& !isTechnicalCozeFailure(row.getError()))
.count();
if (unresolved <= 0) {
return "";
}
return "Coze async workflow returned empty result rows: " + unresolved + "/" + Math.max(expectedRows, cozeRows.size());
}
private void updateCozeStateRunning(TaskScopeStateEntity state, String error) {
int updated = taskScopeStateMapper.update(null, new LambdaUpdateWrapper<TaskScopeStateEntity>()
.eq(TaskScopeStateEntity::getId, state.getId())

View File

@@ -36,7 +36,7 @@ AIIMAGE_APPEARANCE_PATENT_STALE_TIMEOUT_MINUTES=20
AIIMAGE_SIMILAR_ASIN_COZE_BASE_URL=https://api.coze.cn
AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_PATH=/v1/workflow/run
AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID=7632683471312355338
AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID=7635328462404583478
AIIMAGE_SIMILAR_ASIN_COZE_TOKEN=
AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE=50
AIIMAGE_SIMILAR_ASIN_COZE_READ_TIMEOUT_MILLIS=60000

View File

@@ -157,7 +157,7 @@ aiimage:
similar-asin:
coze-base-url: ${AIIMAGE_SIMILAR_ASIN_COZE_BASE_URL:https://api.coze.cn}
coze-workflow-path: ${AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_PATH:/v1/workflow/run}
coze-workflow-id: ${AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID:7632683471312355338}
coze-workflow-id: ${AIIMAGE_SIMILAR_ASIN_COZE_WORKFLOW_ID:7635328462404583478}
coze-token: ${AIIMAGE_SIMILAR_ASIN_COZE_TOKEN:Bearer sat_gsInckdsN0qqhnZSK7K4Zevw6neJCHmYHyIbm85oGgJIGEmSdkH5OFctSKLsAHvT}
coze-batch-size: ${AIIMAGE_SIMILAR_ASIN_COZE_BATCH_SIZE:50}
coze-connect-timeout-millis: ${AIIMAGE_SIMILAR_ASIN_COZE_CONNECT_TIMEOUT_MILLIS:10000}