| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423 |
- /**
- * Sorftime API 数据采集定时任务
- *
- * 功能:
- * - 每月1号凌晨获取类目热销产品BRS 400
- * - 每周一凌晨1点通过类目反查关键词
- *
- * 使用方式:
- * import { sorftimeScheduler } from './modules/sorftime-api-schedule.ts';
- * await sorftimeScheduler.start();
- */
- // node-cron 将在需要时动态加载,避免编译时依赖问题
- import nodeCron from 'npm:node-cron';
- import { relayClient } from '../../../src/relay/relay-client.ts';
- function unwrapRelayData(result) {
- let data = result?.data ?? result?.Data ?? result;
- if (data && typeof data === 'object' && !Array.isArray(data)) {
- data = data.data ?? data.Data ?? data;
- }
- return data;
- }
- function relaySucceeded(result) {
- const code = Number(result?.code ?? result?.Code ?? 200);
- return code === 0 || code === 200;
- }
- function siteCode(marketplaceId, domain) {
- const byMarketplace = {
- ATVPDKIKX0DER: 'us', A2EUQ1WTGCTBG2: 'ca', A1AM78C64UM0Y8: 'mx', A2Q3Y263D00KWC: 'br',
- A1F83G8C2ARO7P: 'uk', A1PA6795UKMFR9: 'de', A13V1IB3VIYZZH: 'fr', A1RKKUPIHCS9HS: 'es',
- APJ6JRA9NG5V4: 'it', A1VC38T7YXB528: 'jp', A39IBJ37TRP1C6: 'au', A2VIGQ35RCS4UG: 'ae',
- A17E79C6D8DWNP: 'sa', A21TJRUUN4KGV: 'in',
- };
- return byMarketplace[marketplaceId] || String(domain || '');
- }
- async function upsertSorftimeProduct(raw, shopId, shopName, marketplaceId, domain) {
- const Parse = globalThis.Parse;
- const asin = String(raw?.Asin || raw?.ASIN || raw?.asin || '').trim().toUpperCase();
- if (!asin) return false;
- const shop = Parse.Object.extend('Shop').createWithoutData(shopId);
- const common = {
- asin,
- parentAsin: raw.ParentAsin || raw.ParentASIN || '',
- title: raw.Title || raw.title || '',
- photo: raw.Photo || raw.photo || '',
- imageUrl: Array.isArray(raw.Photo) ? raw.Photo[0] : (raw.Photo || ''),
- price: Number(raw.Price || 0),
- salesPrice: Number(raw.SalesPrice || raw.Price || 0),
- brand: raw.Brand || '',
- sellerId: raw.BuyboxSellerId || '',
- ratings: Number(raw.Ratings || 0),
- rating: Number(raw.Ratings || 0),
- ratingsCount: Number(raw.RatingsCount || 0),
- category: raw.Category || '',
- bsrCategory: raw.BsrCategory || [],
- rank: Number(raw.Rank || 0),
- listingSalesVolumeOfMonth: Number(raw.ListingSalesVolumeOfMonth || 0),
- ListingSalesVolumeOfMonth: Number(raw.ListingSalesVolumeOfMonth || 0),
- listingSalesOfMonth: Number(raw.ListingSalesOfMonth || 0),
- marketplaceId,
- domain: String(domain || ''),
- site: siteCode(marketplaceId, domain),
- storeName: shopName || '',
- shopName: shopName || '',
- shopId,
- shop,
- source: 'sorftime',
- onlineDate: String(raw.OnlineDate || '').slice(0, 10),
- rawData: raw,
- };
- for (const className of ['Product', 'ProductDetail', 'SorftimeProduct']) {
- const query = new Parse.Query(className);
- query.equalTo('asin', asin);
- if (className !== 'SorftimeProduct') query.equalTo('shop', shop);
- let object = null;
- try {
- object = await query.first({ useMasterKey: true });
- } catch (error) {
- const message = String(error?.message || error || '');
- if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
- }
- if (!object) object = new Parse.Object(className);
- for (const [key, value] of Object.entries(common)) {
- if (value !== undefined && value !== null) object.set(key, value);
- }
- await object.save(null, { useMasterKey: true });
- }
- return true;
- }
- async function upsertSorftimeReview(raw, asin, shopId) {
- const Parse = globalThis.Parse;
- const reviewId = String(raw?.ReviewId || raw?.Id || raw?.ID || raw?.id || '').trim();
- const uniqueId = reviewId || `${asin}:${raw?.ReviewsDate || raw?.ReviewDate || raw?.Date || ''}:${raw?.ProfileName || raw?.Author || ''}`;
- const query = new Parse.Query('SorftimeReviews');
- query.equalTo('reviewId', uniqueId);
- let object = null;
- try {
- object = await query.first({ useMasterKey: true });
- } catch (error) {
- const message = String(error?.message || error || '');
- if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
- }
- if (!object) object = new Parse.Object('SorftimeReviews');
- object.set('reviewId', uniqueId);
- object.set('asin', asin);
- object.set('star', Number(raw?.Star || raw?.Rating || raw?.rating || 0));
- object.set('rating', Number(raw?.Star || raw?.Rating || raw?.rating || 0));
- object.set('title', raw?.Title || raw?.ReviewTitle || '');
- object.set('content', raw?.Content || raw?.ReviewContent || raw?.Body || '');
- object.set('reviewDate', raw?.ReviewsDate || raw?.ReviewDate || raw?.Date || '');
- object.set('author', raw?.ConsumerName || raw?.ProfileName || raw?.Author || '');
- object.set('verified', Boolean(raw?.IsVP || raw?.OnlyPurchase || raw?.VerifiedPurchase));
- object.set('shop', Parse.Object.extend('Shop').createWithoutData(shopId));
- object.set('rawData', raw);
- await object.save(null, { useMasterKey: true });
- }
- class SorftimeScheduler {
- // 私有属性(先声明)
- #isRunning = false;
- #lastRunTime = null;
- #categoryProductsCronTask = null;
- #keywordsCronTask = null;
- #marketTrendCronTask = null;
- #shopProductsCronTask = null;
- #productDetailCronTask = null;
- #cloudFunctionCronTask = null;
- #reviewsCronTask = null;
- #cronEnabled = true;
- // ========== 私有工具方法(最优先声明,避免调用时未定义) ==========
- /**
- * 延迟函数
- * @private
- * @param {number} ms - 延迟毫秒数
- * @returns {Promise<void>} 延迟 Promise
- */
- #delay(ms) {
- return new Promise(resolve => setTimeout(resolve, ms));
- }
- /**
- * 获取所有叶子节点类目
- * @private
- * @returns {Promise<any[]>} 叶子节点类目列表
- */
- async #getLeafCategories() {
- const Parse = globalThis.Parse;
- const query = new Parse.Query('SelfCategory');
- query.equalTo('isLeaf', true);
-
- query.limit(10000);
- return await query.find({ useMasterKey: true });
- }
- /**
- * 根据 shopId 或 nodeIds 筛选类目
- * - nodeIds 优先:直接按 nodeId 数组过滤 SelfCategory
- * - shopId:读取店铺的 nodeIds 字段后过滤 SelfCategory
- * - 两者均未传:返回全部叶子节点类目
- * @private
- * @param {string} shopId - 可选
- * @param {string[]} nodeIds - 可选
- * @returns {Promise<any[]>} 类目列表
- */
- async #getLeafCategoriesByFilter(shopId, nodeIds) {
- const Parse = globalThis.Parse;
- if (nodeIds && nodeIds.length > 0) {
- const query = new Parse.Query('SelfCategory');
- query.containedIn('nodeId', nodeIds);
- query.limit(10000);
- return await query.find({ useMasterKey: true });
- }
- if (shopId) {
- const shopQuery = new Parse.Query('Shop');
- const shop = await shopQuery.get(shopId, { useMasterKey: true });
- const shopNodeIds = shop.get('nodeIds') || [];
- if (shopNodeIds.length > 0) {
- const catQuery = new Parse.Query('SelfCategory');
- catQuery.containedIn('nodeId', shopNodeIds);
- catQuery.limit(10000);
- return await catQuery.find({ useMasterKey: true });
- }
- }
- return await this.#getLeafCategories();
- }
- /**
- * 获取活跃的 Amazon 店铺列表
- * @private
- * @param {string} shopId - 可选,传入则只返回该店铺
- * @returns {Promise<any[]>} 店铺列表
- */
- async #getActiveAmazonShops(shopId) {
- const Parse = globalThis.Parse;
- const query = new Parse.Query('Shop');
- query.equalTo('platform', 'amazon');
- query.equalTo('status', 'active');
- if (shopId) query.equalTo('objectId', shopId);
- query.limit(1000);
- return await query.find({ useMasterKey: true });
- }
- /**
- * 记录执行日志
- * @private
- * @param {any} logData - 日志数据对象
- * @returns {Promise<void>}
- */
- async #logExecution(logData) {
- try {
- const Parse = globalThis.Parse;
- const TaskLog = Parse.Object.extend('TaskExecutionLog');
- const log = new TaskLog();
- log.set('taskName', logData.taskName);
- log.set('startTime', logData.startTime);
- log.set('endTime', logData.endTime);
- log.set('duration', logData.duration);
- log.set('successCount', logData.successCount);
- log.set('errorCount', logData.errorCount);
- log.set('errors', logData.errors);
- log.set('status', logData.status);
- await log.save(null, { useMasterKey: true });
- console.log('[Sorftime Scheduler] 执行日志已保存');
- } catch (error) {
- console.error('[Sorftime Scheduler] 保存执行日志失败:', error.message);
- }
- }
- /**
- * 采集类目热销产品(每月1号凌晨执行)
- * @private
- * @param {string} nodeId - 类目节点ID
- * @param {number} domain - 站点域名代码
- * @returns {Promise<void>}
- */
- async #collectCategoryProducts(nodeId, domain, progress) {
- console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的热销产品 (domain: ${domain})`);
-
- // 固定的云函数ID
- const functionId = 'ZmYNPsoX9X';
-
- // 每页100条,共需要请求4次获取400条数据
- for (let page = 1; page <= 4; page++) {
- try {
- progress?.({ type: 'info', message: `类目 ${nodeId} 第 ${page}/4 页热销产品采集中...` });
- const requestBody = {
- path: "/api/CategoryProducts",
- method: "POST",
- query: { domain },
- body: {
- NodeId: nodeId,
- Page: page,
- Range: 400
- },
- functionId
- };
-
- const result = await relayClient.forwardSorftime({
- path: requestBody.path,
- method: requestBody.method,
- query: requestBody.query,
- body: requestBody.body
- });
-
- if (!relaySucceeded(result)) {
- throw new Error(result.message || `第${page}页数据采集失败`);
- }
- console.log(`[Sorftime Scheduler] 类目 ${nodeId} 第${page}页数据采集完成`);
- // 请求间隔,避免频率过高
- await this.#delay(2000);
- } catch (error) {
- console.error(`[Sorftime Scheduler] 类目 ${nodeId} 第${page}页采集失败:`, error.message);
- throw error;
- }
- }
- }
- /**
- * 采集类目关键词(每周一凌晨1点执行)
- * @private
- * @param {string} nodeId - 类目节点ID
- * @param {number} domain - 站点域名代码
- * @returns {Promise<void>}
- */
- async #collectCategoryKeywords(nodeId, domain, progress) {
- console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的关键词 (domain: ${domain})`);
-
- // 固定的云函数ID
- const functionId = 'KNN4L19aoi';
-
- try {
- const requestBody = {
- path: "/api/CategoryRequestKeyword",
- method: "POST",
- query: { domain },
- body: {
- Nodeid: nodeId,
- PageIndex: 1,
- PageSize: 100
- },
- functionId
- };
- const result = await relayClient.forwardSorftime({
- path: requestBody.path,
- method: requestBody.method,
- query: requestBody.query,
- body: requestBody.body
- });
-
- if (!relaySucceeded(result)) {
- throw new Error(result.message || '关键词采集失败');
- }
- console.log(`[Sorftime Scheduler] 类目 ${nodeId} 关键词采集完成`);
- } catch (error) {
- console.error(`[Sorftime Scheduler] 类目 ${nodeId} 关键词采集失败:`, error.message);
- throw error;
- }
- }
- /**
- * 通过 SellerId 分页采集单个店铺的产品数据
- * @private
- * @param {string} shopId - 店铺 Parse objectId
- * @param {string} sellerId - 亚马逊 SellerId
- * @param {number} domain - 站点代码(1=US)
- * @param {Function} progress - 可选进度回调
- * @returns {Promise<{processed: number, pages: number}>}
- */
- async #doCollectShopProducts(shopId, shopName, sellerId, domain, marketplaceId, progress) {
- const functionId = 'GSjAsvw9FK';
- let page = 1;
- let hasMore = true;
- let totalProcessed = 0;
- progress?.({ type: 'info', message: `SellerId: ${sellerId},开始分页采集产品数据...` });
- console.log(`[Sorftime Scheduler] 店铺 ${shopId} SellerId: ${sellerId} 开始采集`);
- while (hasMore) {
- progress?.({ type: 'info', message: `正在采集第 ${page} 页...` });
- console.log(`[Sorftime Scheduler] 店铺 ${shopId} 采集第 ${page} 页`);
- const result = await relayClient.forwardSorftime({
- path: '/api/ProductQuery',
- method: 'POST',
- query: { domain },
- body: {
- Query: 1,
- QueryType: 5,
- Pattern: sellerId,
- Page: page
- }
- });
- if (!relaySucceeded(result)) {
- throw new Error(`第 ${page} 页采集失败: ${result.message || result.Message || JSON.stringify(result)}`);
- }
- const payload = unwrapRelayData(result);
- const items = payload?.Items || payload?.Products || (Array.isArray(payload) ? payload : []);
- const count = Array.isArray(items) ? items.length : 0;
- let saved = 0;
- for (const item of items) {
- if (await upsertSorftimeProduct(item, shopId, shopName, item?.MarketplaceId || marketplaceId || '', domain)) saved++;
- }
- totalProcessed += count;
- progress?.({ type: 'info', message: `第 ${page} 页获取 ${count} 条产品,入库 ${saved} 条`, processed: totalProcessed });
- console.log(`[Sorftime Scheduler] 店铺 ${shopId} 第 ${page} 页获取 ${count} 条产品,入库 ${saved} 条,累计 ${totalProcessed} 条`);
- const pageCount = Number(payload?.PageCount || 0);
- if (count === 0 || (pageCount > 0 && page >= pageCount) || page >= 100) {
- hasMore = false;
- } else {
- page++;
- await this.#delay(2000);
- }
- }
- return { processed: totalProcessed, pages: page };
- }
- /**
- * 采集类目市场趋势(每月1号凌晨1点执行)
- * @private
- * @param {string} nodeId - 类目节点ID
- * @param {number} domain - 站点域名代码
- * @returns {Promise<void>}
- */
- async #collectCategoryMarketTrend(nodeId, domain, progress) {
- console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的市场趋势 (domain: ${domain})`);
- const functionId = 'OCOKi93GjS';
- const trendIndexes = [0, 1, 2, 3, 4, 5];
- for (const TrendIndex of trendIndexes) {
- try {
- progress?.({ type: 'info', message: `类目 ${nodeId} 趋势类型 ${TrendIndex}/5 采集中...` });
- const requestBody = {
- path: "/api/CategoryTrend",
- method: "POST",
- query: { domain },
- body: {
- NodeId: nodeId,
- TrendIndex
- },
- functionId
- };
- const result = await relayClient.forwardSorftime({
- path: requestBody.path,
- method: requestBody.method,
- query: requestBody.query,
- body: requestBody.body
- });
- if (!relaySucceeded(result)) {
- throw new Error(result.message || `市场趋势类型 ${TrendIndex} 采集失败`);
- }
- console.log(`[Sorftime Scheduler] 类目 ${nodeId} 市场趋势类型 ${TrendIndex} 采集完成`);
- await this.#delay(1000);
- } catch (error) {
- console.error(`[Sorftime Scheduler] 类目 ${nodeId} 市场趋势类型 ${TrendIndex} 采集失败:`, error.message);
- throw error;
- }
- }
- }
- /**
- * 执行类目市场趋势采集任务(每月1号凌晨1点)
- * @private
- * @returns {Promise<any>} 执行结果对象
- */
- async #executeCategoryMarketTrendCollection(shopId, nodeIds, progress) {
- if (this.#isRunning) {
- console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
- return {
- success: false,
- message: '任务正在执行中',
- successCount: 0,
- errorCount: 0,
- duration: 0,
- errors: []
- };
- }
- this.#isRunning = true;
- const startTime = Date.now();
- const errors = [];
- let successCount = 0;
- let errorCount = 0;
- console.log('[Sorftime Scheduler] ========================================');
- console.log('[Sorftime Scheduler] 开始执行类目市场趋势采集任务');
- console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
- console.log('[Sorftime Scheduler] ========================================');
- try {
- const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
- if (categories.length === 0) {
- console.log('[Sorftime Scheduler] 没有找到叶子节点类目数据');
- return {
- success: true,
- message: '没有叶子节点类目需要处理',
- successCount: 0,
- errorCount: 0,
- duration: Date.now() - startTime,
- errors: []
- };
- }
- console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个叶子节点类目`);
- for (const category of categories) {
- try {
- const nodeId = category.get('nodeId');
- const categoryName = category.get('name');
- const domain = category.get('domain') || 1;
- console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
- console.log(`[Sorftime Scheduler] 开始处理类目市场趋势: ${categoryName} (${nodeId}, domain: ${domain})`);
- progress?.({ type: 'info', message: `开始采集类目【${categoryName}】市场趋势` });
- await this.#collectCategoryMarketTrend(nodeId, domain, progress);
- successCount++;
- progress?.({ type: 'success', message: `类目【${categoryName}】市场趋势采集完成`, processed: successCount });
- console.log(`[Sorftime Scheduler] 类目 ${categoryName} 市场趋势采集完成`);
- await this.#delay(2000);
- } catch (error) {
- errorCount++;
- const errorMsg = `类目 ${category.get('name')} 市场趋势采集失败: ${error.message}`;
- errors.push(errorMsg);
- progress?.({ type: 'error', message: errorMsg });
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- }
- }
- const duration = Date.now() - startTime;
- await this.#logExecution({
- taskName: 'sorftime-category-trend-monthly',
- startTime: new Date(startTime),
- endTime: new Date(),
- duration,
- successCount,
- errorCount,
- errors,
- status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
- });
- this.#lastRunTime = new Date();
- const message = `类目市场趋势采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
- console.log(`\n[Sorftime Scheduler] ========================================`);
- console.log(`[Sorftime Scheduler] ${message}`);
- console.log(`[Sorftime Scheduler] ========================================\n`);
- return {
- success: errorCount === 0,
- message,
- successCount,
- errorCount,
- duration,
- errors
- };
- } catch (error) {
- const duration = Date.now() - startTime;
- const errorMsg = `类目市场趋势采集任务执行失败: ${error.message}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- return {
- success: false,
- message: errorMsg,
- successCount,
- errorCount: errorCount + 1,
- duration,
- errors: [...errors, errorMsg]
- };
- } finally {
- this.#isRunning = false;
- }
- }
- /**
- * 执行类目热销产品采集任务(每月1号)
- * @private
- * @returns {Promise<any>} 执行结果对象,包含成功数、失败数、耗时等信息
- */
- async #executeCategoryProductsCollection(shopId, nodeIds, progress) {
- if (this.#isRunning) {
- console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
- return {
- success: false,
- message: '任务正在执行中',
- successCount: 0,
- errorCount: 0,
- duration: 0,
- errors: []
- };
- }
- this.#isRunning = true;
- const startTime = Date.now();
- const errors = [];
- let successCount = 0;
- let errorCount = 0;
- console.log('[Sorftime Scheduler] ========================================');
- console.log('[Sorftime Scheduler] 开始执行类目热销产品采集任务');
- console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
- console.log('[Sorftime Scheduler] ========================================');
- try {
- // 获取所有类目
- const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
- if (categories.length === 0) {
- console.log('[Sorftime Scheduler] 没有找到类目数据');
- return {
- success: true,
- message: '没有类目需要处理',
- successCount: 0,
- errorCount: 0,
- duration: Date.now() - startTime,
- errors: []
- };
- }
- console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个类目`);
- // 循环处理每个类目
- for (const category of categories) {
- try {
- const nodeId = category.get('nodeId');
- const categoryName = category.get('name');
- const domain = category.get('domain') || 1; // 从类目中获取domain,默认为1(美国站)
- console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
- console.log(`[Sorftime Scheduler] 开始处理类目: ${categoryName} (${nodeId}, domain: ${domain})`);
- progress?.({ type: 'info', message: `开始采集类目【${categoryName}】热销产品` });
- await this.#collectCategoryProducts(nodeId, domain, progress);
- successCount++;
- progress?.({ type: 'success', message: `类目【${categoryName}】热销产品采集完成`, processed: successCount });
- console.log(`[Sorftime Scheduler] 类目 ${categoryName} 处理完成`);
- await this.#delay(3000);
- } catch (error) {
- errorCount++;
- const errorMsg = `类目 ${category.get('name')} 处理失败: ${error.message}`;
- errors.push(errorMsg);
- progress?.({ type: 'error', message: errorMsg });
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- }
- }
- // 记录执行日志
- const duration = Date.now() - startTime;
- await this.#logExecution({
- taskName: 'sorftime-category-products-monthly',
- startTime: new Date(startTime),
- endTime: new Date(),
- duration,
- successCount,
- errorCount,
- errors,
- status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
- });
- this.#lastRunTime = new Date();
- const message = `类目热销产品采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
-
- console.log(`\n[Sorftime Scheduler] ========================================`);
- console.log(`[Sorftime Scheduler] ${message}`);
- console.log(`[Sorftime Scheduler] ========================================\n`);
- return {
- success: errorCount === 0,
- message,
- successCount,
- errorCount,
- duration,
- errors
- };
- } catch (error) {
- const duration = Date.now() - startTime;
- const errorMsg = `类目热销产品采集任务执行失败: ${error.message}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- return {
- success: false,
- message: errorMsg,
- successCount,
- errorCount: errorCount + 1,
- duration,
- errors: [...errors, errorMsg]
- };
- } finally {
- this.#isRunning = false;
- }
- }
- /**
- * 执行类目关键词采集任务(每周一)
- * @private
- * @returns {Promise<any>} 执行结果对象,包含成功数、失败数、耗时等信息
- */
- async #executeCategoryKeywordsCollection(shopId, nodeIds, progress) {
- if (this.#isRunning) {
- console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
- return {
- success: false,
- message: '任务正在执行中',
- successCount: 0,
- errorCount: 0,
- duration: 0,
- errors: []
- };
- }
- this.#isRunning = true;
- const startTime = Date.now();
- const errors = [];
- let successCount = 0;
- let errorCount = 0;
- console.log('[Sorftime Scheduler] ========================================');
- console.log('[Sorftime Scheduler] 开始执行类目关键词采集任务');
- console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
- console.log('[Sorftime Scheduler] ========================================');
- try {
- // 获取所有叶子节点类目
- const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
- if (categories.length === 0) {
- console.log('[Sorftime Scheduler] 没有找到叶子节点类目数据');
- return {
- success: true,
- message: '没有叶子节点类目需要处理',
- successCount: 0,
- errorCount: 0,
- duration: Date.now() - startTime,
- errors: []
- };
- }
- console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个叶子节点类目`);
- // 循环处理每个类目
- for (const category of categories) {
- try {
- const nodeId = category.get('nodeId');
- const categoryName = category.get('name');
- const domain = category.get('domain') || 1; // 从类目中获取domain,默认为1(美国站)
-
- console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
- console.log(`[Sorftime Scheduler] 开始处理类目: ${categoryName} (${nodeId}, domain: ${domain})`);
- progress?.({ type: 'info', message: `开始采集类目【${categoryName}】关键词` });
- await this.#collectCategoryKeywords(nodeId, domain, progress);
- successCount++;
-
- progress?.({ type: 'success', message: `类目【${categoryName}】关键词采集完成`, processed: successCount });
- console.log(`[Sorftime Scheduler] 类目 ${categoryName} 关键词采集完成`);
- await this.#delay(2000);
- } catch (error) {
- errorCount++;
- const errorMsg = `类目 ${category.get('name')} 关键词采集失败: ${error.message}`;
- errors.push(errorMsg);
- progress?.({ type: 'error', message: errorMsg });
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- }
- }
- // 记录执行日志
- const duration = Date.now() - startTime;
- await this.#logExecution({
- taskName: 'sorftime-category-keywords-weekly',
- startTime: new Date(startTime),
- endTime: new Date(),
- duration,
- successCount,
- errorCount,
- errors,
- status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
- });
- this.#lastRunTime = new Date();
- const message = `类目关键词采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
-
- console.log(`\n[Sorftime Scheduler] ========================================`);
- console.log(`[Sorftime Scheduler] ${message}`);
- console.log(`[Sorftime Scheduler] ========================================\n`);
- return {
- success: errorCount === 0,
- message,
- successCount,
- errorCount,
- duration,
- errors
- };
- } catch (error) {
- const duration = Date.now() - startTime;
- const errorMsg = `类目关键词采集任务执行失败: ${error.message}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- return {
- success: false,
- message: errorMsg,
- successCount,
- errorCount: errorCount + 1,
- duration,
- errors: [...errors, errorMsg]
- };
- } finally {
- this.#isRunning = false;
- }
- }
- /**
- * 执行云函数调用任务(每两小时)
- * 调用指定的云函数进行数据处理
- * @private
- * @returns {Promise<any>} 执行结果对象
- */
- async #executeCloudFunctionInvocation() {
- console.log('[Sorftime Scheduler] ========================================');
- console.log('[Sorftime Scheduler] 开始执行云函数调用任务');
- console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
- console.log('[Sorftime Scheduler] ========================================');
- const startTime = Date.now();
- try {
- const response = await fetch('http://localhost:3000/api/functions', {
- method: 'POST',
- headers: {
- 'Content-Type': 'application/json'
- },
- body: JSON.stringify({
- id: "Z4z6SB4o9e"
- })
- });
- const duration = Date.now() - startTime;
- if (!response.ok) {
- const errorMsg = `云函数调用失败,状态码: ${response.status}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- return {
- success: false,
- message: errorMsg,
- duration,
- errors: [errorMsg]
- };
- }
- const result = await response.json();
- const message = `云函数调用完成,耗时 ${Math.round(duration / 1000)}秒`;
- console.log(`\n[Sorftime Scheduler] ========================================`);
- console.log(`[Sorftime Scheduler] ${message}`);
- console.log(`[Sorftime Scheduler] ========================================\n`);
- return {
- success: true,
- message,
- duration,
- data: result,
- errors: []
- };
- } catch (error) {
- const duration = Date.now() - startTime;
- const errorMsg = `云函数调用任务执行失败: ${error.message}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- return {
- success: false,
- message: errorMsg,
- duration,
- errors: [errorMsg]
- };
- }
- }
- /**
- * 执行产品评论采集任务(每天凌晨1点)
- * 根据 Product 表中的 ASIN,采集对应的评论信息
- * @private
- * @returns {Promise<any>} 执行结果对象
- */
- async #executeProductReviewsCollection(progress) {
- console.log('[Sorftime Scheduler] ========================================');
- console.log('[Sorftime Scheduler] 开始执行产品评论采集任务');
- console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
- console.log('[Sorftime Scheduler] ========================================');
- const startTime = Date.now();
- let successCount = 0;
- let errorCount = 0;
- const errors = [];
- progress?.({ type: 'info', message: '开始执行产品评论采集任务' });
- try {
- const Parse = globalThis.Parse;
- // 查询所有 Product 记录
- const productQuery = new Parse.Query('Product');
- productQuery.limit(10000);
- const products = await productQuery.find({ useMasterKey: true });
- if (products.length === 0) {
- console.log('[Sorftime Scheduler] 没有找到 Product 数据');
- progress?.({ type: 'info', message: '没有 Product 数据需要处理' });
- return {
- success: true,
- message: '没有 Product 数据需要处理',
- successCount: 0,
- errorCount: 0,
- duration: Date.now() - startTime,
- errors: []
- };
- }
- console.log(`[Sorftime Scheduler] 发现 ${products.length} 个 Product 记录`);
- progress?.({ type: 'info', message: `共发现 ${products.length} 个产品需要检查评论` });
- // 遍历每个 Product
- for (const product of products) {
- try {
- const asin = product.get('asin');
- if (!asin) {
- console.log('[Sorftime Scheduler] 跳过无 ASIN 的 Product');
- continue;
- }
- progress?.({ type: 'info', message: `开始处理 ASIN: ${asin}` });
- // 获取 Product 关联的 Shop 信息
- const shopRelation = product.get('shop');
- if (!shopRelation) {
- console.log(`[Sorftime Scheduler] Product 无关联 Shop,ASIN: ${asin}`);
- continue;
- }
- const shop = await shopRelation.fetch({ useMasterKey: true });
- const shopDomain = shop.get('domain') || 1;
- // 查询 SorftimeReviews 中该 ASIN 最新的 updateAt 时间
- let queryStartDt = '2025-01-01';
- const reviewQuery = new Parse.Query('SorftimeReviews');
- reviewQuery.equalTo('asin', asin);
- reviewQuery.descending('updatedAt');
- reviewQuery.limit(1);
- let latestReview = null;
- try {
- latestReview = await reviewQuery.first({ useMasterKey: true });
- } catch (error) {
- const message = String(error?.message || error || '');
- if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
- }
- if (latestReview) {
- const updateAt = latestReview.get('updatedAt');
- if (updateAt) {
- // 格式化为 YYYY-MM-DD
- const date = new Date(updateAt);
- queryStartDt = date.toISOString().split('T')[0];
- }
- }
- console.log(`[Sorftime Scheduler] 开始采集评论,ASIN: ${asin}, 查询起始日期: ${queryStartDt}`);
- // 调用 Sorftime 接口获取评论信息
- let result;
- try {
- result = await relayClient.forwardSorftime({
- path: '/api/ProductReviewsQuery',
- method: 'POST',
- query: {
- domain: shopDomain,
- shopId: shop.id
- },
- body: {
- ASIN: asin,
- PageIndex: 1,
- OnlyPurchase: 1,
- Star: '1,2,3,4,5',
- QueryStartDt: queryStartDt
- }
- });
- } catch (error) {
- const errorMsg = `请求评论数据失败,ASIN: ${asin}, ${error.message || error}`;
- errors.push(errorMsg);
- errorCount++;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- continue;
- }
- const payload = unwrapRelayData(result);
- const reviews = payload?.Reviews || payload?.Items || payload?.List || (Array.isArray(payload) ? payload : []);
- if (relaySucceeded(result)) {
- for (const review of reviews) await upsertSorftimeReview(review, asin, shop.id);
- successCount++;
- const statusText = reviews.length ? `入库 ${reviews.length} 条` : '无新增评论';
- console.log(`[Sorftime Scheduler] 评论数据采集成功,ASIN: ${asin},${statusText}`);
- progress?.({ type: 'success', message: `评论数据采集成功,ASIN: ${asin},${statusText}`, processed: successCount });
- } else {
- const errorMsg = `评论数据获取失败,ASIN: ${asin}`;
- errors.push(errorMsg);
- errorCount++;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- progress?.({ type: 'error', message: errorMsg });
- }
- // 延迟 1 秒,避免请求过于频繁
- await this.#delay(500);
- } catch (error) {
- errorCount++;
- const errorMsg = `处理 Product 失败: ${error.message}`;
- errors.push(errorMsg);
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- progress?.({ type: 'error', message: errorMsg });
- }
- }
- const duration = Date.now() - startTime;
- const message = `产品评论采集完成: 成功 ${successCount}/${products.length} 个产品,耗时 ${Math.round(duration / 1000)}秒`;
- console.log(`\n[Sorftime Scheduler] ========================================`);
- console.log(`[Sorftime Scheduler] ${message}`);
- console.log(`[Sorftime Scheduler] ========================================\n`);
- progress?.({ type: errorCount === 0 ? 'success' : 'warning', message });
- return {
- success: errorCount === 0,
- message,
- successCount,
- errorCount,
- duration,
- errors
- };
- } catch (error) {
- const duration = Date.now() - startTime;
- const errorMsg = `产品评论采集任务执行失败: ${error.message}`;
- console.error(`[Sorftime Scheduler] ${errorMsg}`);
- progress?.({ type: 'error', message: errorMsg });
- return {
- success: false,
- message: errorMsg,
- successCount: 0,
- errorCount: 1,
- duration,
- errors: [errorMsg]
- };
- }
- }
- // ========== 公共方法(后声明,因为依赖前面的私有方法) ==========
- /**
- * 构造函数
- * 初始化 Sorftime API 数据采集调度器
- */
- constructor() {
- console.log('[Sorftime Scheduler] 初始化 Sorftime API 数据采集调度器');
- }
- /**
- * 启动定时任务
- * @returns {Promise<void>}
- */
- async start() {
- if (this.#categoryProductsCronTask || this.#keywordsCronTask || this.#marketTrendCronTask || this.#shopProductsCronTask || this.#productDetailCronTask || this.#cloudFunctionCronTask || this.#reviewsCronTask) {
- console.log('[Sorftime Scheduler] 定时任务已在运行中');
- return;
- }
- // 每月1号凌晨执行类目热销产品采集
- this.#categoryProductsCronTask = nodeCron.schedule('0 0 1 * *', async () => {
- await this.#executeCategoryProductsCollection(undefined, undefined, undefined);
- }, {
- timezone: "Asia/Shanghai"
- });
- console.log('[Sorftime Scheduler] 类目热销产品定时任务已启动,将在每月1号凌晨执行');
- // 每月1号凌晨1点执行类目市场趋势采集
- this.#marketTrendCronTask = nodeCron.schedule('0 1 1 * *', async () => {
- await this.#executeCategoryMarketTrendCollection(undefined, undefined, undefined);
- }, {
- timezone: "Asia/Shanghai"
- });
- console.log('[Sorftime Scheduler] 类目市场趋势定时任务已启动,将在每月1号凌晨1点执行');
- // 每周一凌晨1点执行类目关键词采集
- this.#keywordsCronTask = nodeCron.schedule('0 1 * * 1', async () => {
- await this.#executeCategoryKeywordsCollection(undefined, undefined, undefined);
- }, {
- timezone: "Asia/Shanghai"
- });
- console.log('[Sorftime Scheduler] 类目关键词定时任务已启动,将在每周一凌晨1点执行');
- if (process.env.ENABLE_SORFTIME_REVIEW_SCHEDULES === 'true') {
- this.#reviewsCronTask = nodeCron.schedule('0 1 * * *', async () => {
- console.log('[Sorftime Scheduler] ===== 凌晨1点:开始产品评论采集 =====');
- try {
- await this.#executeProductReviewsCollection(undefined);
- } catch (e) {
- console.error('[Sorftime Scheduler] 产品评论采集失败:', e.message);
- }
- console.log('[Sorftime Scheduler] ===== 产品评论采集完成 =====');
- }, { timezone: 'Asia/Shanghai' });
- console.log('[Sorftime Scheduler] 产品评论采集定时任务已启动,将在每天凌晨1点执行');
- } else {
- console.log('[Sorftime Scheduler] 产品评论采集定时任务已禁用');
- }
- // 每天凌晨2点采集店铺产品数据
- this.#shopProductsCronTask = nodeCron.schedule('0 2 * * *', async () => {
- console.log('[Sorftime Scheduler] ===== 凌晨2点:开始店铺产品采集 =====');
- try {
- await this.collectShopProducts(undefined, undefined);
- } catch (e) {
- console.error('[Sorftime Scheduler] 店铺产品采集失败:', e.message);
- }
- console.log('[Sorftime Scheduler] ===== 店铺产品采集完成 =====');
- }, { timezone: 'Asia/Shanghai' });
- console.log('[Sorftime Scheduler] 店铺产品采集定时任务已启动,将在每天凌晨2点执行');
- if (process.env.ENABLE_LEGACY_CLOUD_FUNCTION_SCHEDULES === 'true') {
- this.#cloudFunctionCronTask = nodeCron.schedule('0 */2 * * *', async () => {
- try {
- await this.#executeCloudFunctionInvocation();
- } catch (e) {
- console.error('[Sorftime Scheduler] 遗留云函数调用失败:', e.message);
- }
- }, { timezone: 'Asia/Shanghai' });
- } else {
- console.log('[Sorftime Scheduler] 遗留动态云函数定时任务已禁用');
- }
- }
- /**
- * 停止定时任务
- * @returns {void}
- */
- stop() {
- if (this.#categoryProductsCronTask) {
- this.#categoryProductsCronTask.stop();
- this.#categoryProductsCronTask = null;
- console.log('[Sorftime Scheduler] 类目热销产品定时任务已停止');
- }
- if (this.#marketTrendCronTask) {
- this.#marketTrendCronTask.stop();
- this.#marketTrendCronTask = null;
- console.log('[Sorftime Scheduler] 类目市场趋势定时任务已停止');
- }
- if (this.#keywordsCronTask) {
- this.#keywordsCronTask.stop();
- this.#keywordsCronTask = null;
- console.log('[Sorftime Scheduler] 类目关键词定时任务已停止');
- }
- if (this.#shopProductsCronTask) {
- this.#shopProductsCronTask.stop();
- this.#shopProductsCronTask = null;
- console.log('[Sorftime Scheduler] 店铺产品采集定时任务已停止');
- }
- if (this.#productDetailCronTask) {
- this.#productDetailCronTask.stop();
- this.#productDetailCronTask = null;
- console.log('[Sorftime Scheduler] 产品详情同步定时任务已停止');
- }
- if (this.#cloudFunctionCronTask) {
- this.#cloudFunctionCronTask.stop();
- this.#cloudFunctionCronTask = null;
- console.log('[Sorftime Scheduler] 云函数调用定时任务已停止');
- }
- if (this.#reviewsCronTask) {
- this.#reviewsCronTask.stop();
- this.#reviewsCronTask = null;
- console.log('[Sorftime Scheduler] 产品评论采集定时任务已停止');
- }
- }
- /**
- * 获取调度器状态
- * @returns {object} 调度器状态对象
- */
- getStatus() {
- return {
- isRunning: this.#isRunning,
- lastRunTime: this.#lastRunTime,
- cronEnabled: this.#cronEnabled,
- tasks: {
- categoryProducts: {
- enabled: !!this.#categoryProductsCronTask,
- schedule: '每月1号凌晨'
- },
- marketTrend: {
- enabled: !!this.#marketTrendCronTask,
- schedule: '每月1号凌晨1点'
- },
- keywords: {
- enabled: !!this.#keywordsCronTask,
- schedule: '每周一凌晨1点'
- },
- shopProducts: {
- enabled: !!this.#shopProductsCronTask,
- schedule: '每天凌晨2点'
- },
- productDetails: {
- enabled: !!this.#productDetailCronTask,
- schedule: '每小时'
- },
- cloudFunction: {
- enabled: !!this.#cloudFunctionCronTask,
- schedule: '每两小时'
- },
- productReviews: {
- enabled: !!this.#reviewsCronTask,
- schedule: '每天凌晨1点'
- }
- }
- };
- }
- /**
- * 手动触发类目热销产品采集
- * @returns {Promise<any>} 执行结果对象
- */
- async triggerCategoryProducts(shopId, nodeIds, progress) {
- console.log('[Sorftime Scheduler] 手动触发类目热销产品采集任务');
- return await this.#executeCategoryProductsCollection(shopId, nodeIds, progress);
- }
- /**
- * 手动触发类目关键词采集
- * @returns {Promise<any>} 执行结果对象
- */
- async triggerCategoryKeywords(shopId, nodeIds, progress) {
- console.log('[Sorftime Scheduler] 手动触发类目关键词采集任务');
- return await this.#executeCategoryKeywordsCollection(shopId, nodeIds, progress);
- }
- /**
- * 手动触发类目市场趋势采集
- * @returns {Promise<any>} 执行结果对象
- */
- async triggerCategoryTrends(shopId, nodeIds, progress) {
- console.log('[Sorftime Scheduler] 手动触发类目市场趋势采集任务');
- return await this.#executeCategoryMarketTrendCollection(shopId, nodeIds, progress);
- }
- /**
- * 手动触发产品评论采集
- * @returns {Promise<any>} 执行结果对象
- */
- async triggerProductReviews(progress) {
- console.log('[Sorftime Scheduler] 手动触发产品评论采集任务');
- return await this.#executeProductReviewsCollection(progress);
- }
- /**
- * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势)
- * 前端选中某个类目但数据库无数据时调用
- * @param {string} nodeId - 类目节点ID
- * @param {number} domain - 站点域名代码(默认1=美国)
- * @returns {Promise<any>} 执行结果对象
- */
- async triggerSingleCategory(nodeId, domain = 1, progress) {
- const startTime = Date.now();
- const errors = [];
- let tasks = { products: false, keywords: false, trends: false };
- console.log(`[Sorftime Scheduler] 按需触发单类目采集: nodeId=${nodeId}, domain=${domain}`);
- try {
- // 1. 热销产品
- try {
- progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 热销产品` });
- await this.#collectCategoryProducts(nodeId, domain, progress);
- tasks.products = true;
- progress?.({ type: 'success', message: `类目 ${nodeId} 热销产品采集完成` });
- console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 热销产品采集完成`);
- } catch (error) {
- const errMsg = error?.message || String(error);
- progress?.({ type: 'error', message: `热销产品采集失败: ${errMsg}` });
- errors.push(`热销产品采集失败: ${errMsg}`);
- console.error(`[Sorftime Scheduler] 热销产品采集失败:`, error);
- }
- // 延迟 2 秒
- await new Promise(resolve => setTimeout(resolve, 500));
- // 2. 关键词
- try {
- progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 关键词` });
- await this.#collectCategoryKeywords(nodeId, domain, progress);
- tasks.keywords = true;
- progress?.({ type: 'success', message: `类目 ${nodeId} 关键词采集完成` });
- console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 关键词采集完成`);
- } catch (error) {
- const errMsg = error?.message || String(error);
- progress?.({ type: 'error', message: `关键词采集失败: ${errMsg}` });
- errors.push(`关键词采集失败: ${errMsg}`);
- console.error(`[Sorftime Scheduler] 关键词采集失败:`, error);
- }
- // 延迟 2 秒
- await new Promise(resolve => setTimeout(resolve, 500));
- // 3. 市场趋势
- try {
- progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 市场趋势` });
- await this.#collectCategoryMarketTrend(nodeId, domain, progress);
- tasks.trends = true;
- progress?.({ type: 'success', message: `类目 ${nodeId} 市场趋势采集完成` });
- console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 市场趋势采集完成`);
- } catch (error) {
- const errMsg = error?.message || String(error);
- progress?.({ type: 'error', message: `市场趋势采集失败: ${errMsg}` });
- errors.push(`市场趋势采集失败: ${errMsg}`);
- console.error(`[Sorftime Scheduler] 市场趋势采集失败:`, error);
- }
- const duration = Date.now() - startTime;
- const allSuccess = errors.length === 0;
- const successCount = Object.values(tasks).filter(Boolean).length;
- // 异步记录执行日志,不阻塞返回
- if (this.#logExecution) {
- this.#logExecution({
- taskName: `trigger-single-category-${nodeId}`,
- startTime: new Date(startTime),
- endTime: new Date(),
- duration,
- successCount,
- errorCount: errors.length,
- errors,
- status: allSuccess ? 'success' : 'partial_success'
- }).catch(err => console.error('[Sorftime Scheduler] Log execution error:', err));
- }
- const result = {
- success: allSuccess,
- message: allSuccess
- ? `类目 ${nodeId} 全量采集完成,耗时 ${Math.round(duration / 1000)}秒`
- : `类目 ${nodeId} 部分采集失败(成功 ${successCount}/3 个任务)`,
- nodeId,
- domain,
- tasks,
- duration,
- errors: errors.length > 0 ? errors : undefined
- };
- console.log(`[Sorftime Scheduler] 采集结果:`, result);
- return result;
- } catch (err) {
- const duration = Date.now() - startTime;
- const errMsg = err?.message || String(err);
- console.error(`[Sorftime Scheduler] 采集异常:`, err);
- return {
- success: false,
- message: `类目 ${nodeId} 采集异常: ${errMsg}`,
- nodeId,
- domain,
- tasks,
- duration,
- errors: [errMsg]
- };
- }
- }
- /**
- * 采集所有活跃 Amazon 店铺(或指定店铺)的产品数据
- * 通过 SellerId (QueryType=5) 分页请求 /api/ProductQuery
- * @param {string} shopId - 可选,不传则采集全部活跃 Amazon 店铺
- * @param {Function} progress - 可选进度回调,适合 SSE 实时推送
- * @returns {Promise<any>} 执行结果
- */
- async collectShopProducts(shopId, progress) {
- const shops = await this.#getActiveAmazonShops(shopId);
- if (shops.length === 0) {
- const msg = shopId ? `未找到店铺 ${shopId} 或店铺未激活` : '没有找到活跃的 Amazon 店铺';
- progress?.({ type: 'warn', message: msg });
- console.warn(`[Sorftime Scheduler] ${msg}`);
- return { success: true, message: msg, processed: 0, errorCount: 0, errors: [] };
- }
- let totalProcessed = 0;
- let errorCount = 0;
- const errors = [];
- for (const shop of shops) {
- const sid = shop.id;
- const shopName = shop.get('name');
- const config = shop.get('config');
- const domain = Number(shop.get('domain') || 1);
- const marketplaceId = shop.get('marketplaceId') || '';
- if (!config || !config.SpApiConfig || !config.SpApiConfig.sellerID) {
- const msg = `店铺【${shopName}】(${sid}) 缺少 SpApiConfig.sellerID 配置,跳过`;
- progress?.({ type: 'error', message: msg });
- console.error(`[Sorftime Scheduler] ${msg}`);
- errors.push(msg);
- errorCount++;
- continue;
- }
- if (config.SpApiConfig.listingEnabled === false) {
- progress?.({ type: 'warn', message: `店铺【${shopName}】未配置区域 Seller ID,跳过 Sorftime 店铺产品采集` });
- continue;
- }
- const sellerId = config.SpApiConfig.sellerID;
- progress?.({ type: 'info', message: `开始采集店铺【${shopName}】的产品数据 (SellerId: ${sellerId})` });
- console.log(`[Sorftime Scheduler] 开始采集店铺 ${shopName} (${sid}) 产品数据`);
- try {
- const result = await this.#doCollectShopProducts(sid, shopName, sellerId, domain, marketplaceId, progress);
- progress?.({ type: 'success', message: `店铺【${shopName}】产品采集完成: 共 ${result.pages} 页 ${result.processed} 条`, processed: result.processed });
- console.log(`[Sorftime Scheduler] 店铺 ${shopName} 产品采集完成: ${result.processed} 条`);
- totalProcessed += result.processed;
- await this.#delay(2000);
- } catch (error) {
- const msg = `店铺【${shopName}】产品采集失败: ${error.message}`;
- progress?.({ type: 'error', message: msg });
- console.error(`[Sorftime Scheduler] ${msg}`);
- errors.push(msg);
- errorCount++;
- }
- }
- const message = `店铺产品采集完成: 成功 ${shops.length - errorCount}/${shops.length} 个店铺,共 ${totalProcessed} 条记录`;
- return { success: errorCount === 0, message, processed: totalProcessed, errorCount, errors };
- }
- }
- // 导出单例
- export const sorftimeScheduler = new SorftimeScheduler();
|