sorftime-api-schedule.ts 51 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423
  1. /**
  2. * Sorftime API 数据采集定时任务
  3. *
  4. * 功能:
  5. * - 每月1号凌晨获取类目热销产品BRS 400
  6. * - 每周一凌晨1点通过类目反查关键词
  7. *
  8. * 使用方式:
  9. * import { sorftimeScheduler } from './modules/sorftime-api-schedule.ts';
  10. * await sorftimeScheduler.start();
  11. */
  12. // node-cron 将在需要时动态加载,避免编译时依赖问题
  13. import nodeCron from 'npm:node-cron';
  14. import { relayClient } from '../../../src/relay/relay-client.ts';
  15. function unwrapRelayData(result) {
  16. let data = result?.data ?? result?.Data ?? result;
  17. if (data && typeof data === 'object' && !Array.isArray(data)) {
  18. data = data.data ?? data.Data ?? data;
  19. }
  20. return data;
  21. }
  22. function relaySucceeded(result) {
  23. const code = Number(result?.code ?? result?.Code ?? 200);
  24. return code === 0 || code === 200;
  25. }
  26. function siteCode(marketplaceId, domain) {
  27. const byMarketplace = {
  28. ATVPDKIKX0DER: 'us', A2EUQ1WTGCTBG2: 'ca', A1AM78C64UM0Y8: 'mx', A2Q3Y263D00KWC: 'br',
  29. A1F83G8C2ARO7P: 'uk', A1PA6795UKMFR9: 'de', A13V1IB3VIYZZH: 'fr', A1RKKUPIHCS9HS: 'es',
  30. APJ6JRA9NG5V4: 'it', A1VC38T7YXB528: 'jp', A39IBJ37TRP1C6: 'au', A2VIGQ35RCS4UG: 'ae',
  31. A17E79C6D8DWNP: 'sa', A21TJRUUN4KGV: 'in',
  32. };
  33. return byMarketplace[marketplaceId] || String(domain || '');
  34. }
  35. async function upsertSorftimeProduct(raw, shopId, shopName, marketplaceId, domain) {
  36. const Parse = globalThis.Parse;
  37. const asin = String(raw?.Asin || raw?.ASIN || raw?.asin || '').trim().toUpperCase();
  38. if (!asin) return false;
  39. const shop = Parse.Object.extend('Shop').createWithoutData(shopId);
  40. const common = {
  41. asin,
  42. parentAsin: raw.ParentAsin || raw.ParentASIN || '',
  43. title: raw.Title || raw.title || '',
  44. photo: raw.Photo || raw.photo || '',
  45. imageUrl: Array.isArray(raw.Photo) ? raw.Photo[0] : (raw.Photo || ''),
  46. price: Number(raw.Price || 0),
  47. salesPrice: Number(raw.SalesPrice || raw.Price || 0),
  48. brand: raw.Brand || '',
  49. sellerId: raw.BuyboxSellerId || '',
  50. ratings: Number(raw.Ratings || 0),
  51. rating: Number(raw.Ratings || 0),
  52. ratingsCount: Number(raw.RatingsCount || 0),
  53. category: raw.Category || '',
  54. bsrCategory: raw.BsrCategory || [],
  55. rank: Number(raw.Rank || 0),
  56. listingSalesVolumeOfMonth: Number(raw.ListingSalesVolumeOfMonth || 0),
  57. ListingSalesVolumeOfMonth: Number(raw.ListingSalesVolumeOfMonth || 0),
  58. listingSalesOfMonth: Number(raw.ListingSalesOfMonth || 0),
  59. marketplaceId,
  60. domain: String(domain || ''),
  61. site: siteCode(marketplaceId, domain),
  62. storeName: shopName || '',
  63. shopName: shopName || '',
  64. shopId,
  65. shop,
  66. source: 'sorftime',
  67. onlineDate: String(raw.OnlineDate || '').slice(0, 10),
  68. rawData: raw,
  69. };
  70. for (const className of ['Product', 'ProductDetail', 'SorftimeProduct']) {
  71. const query = new Parse.Query(className);
  72. query.equalTo('asin', asin);
  73. if (className !== 'SorftimeProduct') query.equalTo('shop', shop);
  74. let object = null;
  75. try {
  76. object = await query.first({ useMasterKey: true });
  77. } catch (error) {
  78. const message = String(error?.message || error || '');
  79. if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
  80. }
  81. if (!object) object = new Parse.Object(className);
  82. for (const [key, value] of Object.entries(common)) {
  83. if (value !== undefined && value !== null) object.set(key, value);
  84. }
  85. await object.save(null, { useMasterKey: true });
  86. }
  87. return true;
  88. }
  89. async function upsertSorftimeReview(raw, asin, shopId) {
  90. const Parse = globalThis.Parse;
  91. const reviewId = String(raw?.ReviewId || raw?.Id || raw?.ID || raw?.id || '').trim();
  92. const uniqueId = reviewId || `${asin}:${raw?.ReviewsDate || raw?.ReviewDate || raw?.Date || ''}:${raw?.ProfileName || raw?.Author || ''}`;
  93. const query = new Parse.Query('SorftimeReviews');
  94. query.equalTo('reviewId', uniqueId);
  95. let object = null;
  96. try {
  97. object = await query.first({ useMasterKey: true });
  98. } catch (error) {
  99. const message = String(error?.message || error || '');
  100. if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
  101. }
  102. if (!object) object = new Parse.Object('SorftimeReviews');
  103. object.set('reviewId', uniqueId);
  104. object.set('asin', asin);
  105. object.set('star', Number(raw?.Star || raw?.Rating || raw?.rating || 0));
  106. object.set('rating', Number(raw?.Star || raw?.Rating || raw?.rating || 0));
  107. object.set('title', raw?.Title || raw?.ReviewTitle || '');
  108. object.set('content', raw?.Content || raw?.ReviewContent || raw?.Body || '');
  109. object.set('reviewDate', raw?.ReviewsDate || raw?.ReviewDate || raw?.Date || '');
  110. object.set('author', raw?.ConsumerName || raw?.ProfileName || raw?.Author || '');
  111. object.set('verified', Boolean(raw?.IsVP || raw?.OnlyPurchase || raw?.VerifiedPurchase));
  112. object.set('shop', Parse.Object.extend('Shop').createWithoutData(shopId));
  113. object.set('rawData', raw);
  114. await object.save(null, { useMasterKey: true });
  115. }
  116. class SorftimeScheduler {
  117. // 私有属性(先声明)
  118. #isRunning = false;
  119. #lastRunTime = null;
  120. #categoryProductsCronTask = null;
  121. #keywordsCronTask = null;
  122. #marketTrendCronTask = null;
  123. #shopProductsCronTask = null;
  124. #productDetailCronTask = null;
  125. #cloudFunctionCronTask = null;
  126. #reviewsCronTask = null;
  127. #cronEnabled = true;
  128. // ========== 私有工具方法(最优先声明,避免调用时未定义) ==========
  129. /**
  130. * 延迟函数
  131. * @private
  132. * @param {number} ms - 延迟毫秒数
  133. * @returns {Promise<void>} 延迟 Promise
  134. */
  135. #delay(ms) {
  136. return new Promise(resolve => setTimeout(resolve, ms));
  137. }
  138. /**
  139. * 获取所有叶子节点类目
  140. * @private
  141. * @returns {Promise<any[]>} 叶子节点类目列表
  142. */
  143. async #getLeafCategories() {
  144. const Parse = globalThis.Parse;
  145. const query = new Parse.Query('SelfCategory');
  146. query.equalTo('isLeaf', true);
  147. query.limit(10000);
  148. return await query.find({ useMasterKey: true });
  149. }
  150. /**
  151. * 根据 shopId 或 nodeIds 筛选类目
  152. * - nodeIds 优先:直接按 nodeId 数组过滤 SelfCategory
  153. * - shopId:读取店铺的 nodeIds 字段后过滤 SelfCategory
  154. * - 两者均未传:返回全部叶子节点类目
  155. * @private
  156. * @param {string} shopId - 可选
  157. * @param {string[]} nodeIds - 可选
  158. * @returns {Promise<any[]>} 类目列表
  159. */
  160. async #getLeafCategoriesByFilter(shopId, nodeIds) {
  161. const Parse = globalThis.Parse;
  162. if (nodeIds && nodeIds.length > 0) {
  163. const query = new Parse.Query('SelfCategory');
  164. query.containedIn('nodeId', nodeIds);
  165. query.limit(10000);
  166. return await query.find({ useMasterKey: true });
  167. }
  168. if (shopId) {
  169. const shopQuery = new Parse.Query('Shop');
  170. const shop = await shopQuery.get(shopId, { useMasterKey: true });
  171. const shopNodeIds = shop.get('nodeIds') || [];
  172. if (shopNodeIds.length > 0) {
  173. const catQuery = new Parse.Query('SelfCategory');
  174. catQuery.containedIn('nodeId', shopNodeIds);
  175. catQuery.limit(10000);
  176. return await catQuery.find({ useMasterKey: true });
  177. }
  178. }
  179. return await this.#getLeafCategories();
  180. }
  181. /**
  182. * 获取活跃的 Amazon 店铺列表
  183. * @private
  184. * @param {string} shopId - 可选,传入则只返回该店铺
  185. * @returns {Promise<any[]>} 店铺列表
  186. */
  187. async #getActiveAmazonShops(shopId) {
  188. const Parse = globalThis.Parse;
  189. const query = new Parse.Query('Shop');
  190. query.equalTo('platform', 'amazon');
  191. query.equalTo('status', 'active');
  192. if (shopId) query.equalTo('objectId', shopId);
  193. query.limit(1000);
  194. return await query.find({ useMasterKey: true });
  195. }
  196. /**
  197. * 记录执行日志
  198. * @private
  199. * @param {any} logData - 日志数据对象
  200. * @returns {Promise<void>}
  201. */
  202. async #logExecution(logData) {
  203. try {
  204. const Parse = globalThis.Parse;
  205. const TaskLog = Parse.Object.extend('TaskExecutionLog');
  206. const log = new TaskLog();
  207. log.set('taskName', logData.taskName);
  208. log.set('startTime', logData.startTime);
  209. log.set('endTime', logData.endTime);
  210. log.set('duration', logData.duration);
  211. log.set('successCount', logData.successCount);
  212. log.set('errorCount', logData.errorCount);
  213. log.set('errors', logData.errors);
  214. log.set('status', logData.status);
  215. await log.save(null, { useMasterKey: true });
  216. console.log('[Sorftime Scheduler] 执行日志已保存');
  217. } catch (error) {
  218. console.error('[Sorftime Scheduler] 保存执行日志失败:', error.message);
  219. }
  220. }
  221. /**
  222. * 采集类目热销产品(每月1号凌晨执行)
  223. * @private
  224. * @param {string} nodeId - 类目节点ID
  225. * @param {number} domain - 站点域名代码
  226. * @returns {Promise<void>}
  227. */
  228. async #collectCategoryProducts(nodeId, domain, progress) {
  229. console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的热销产品 (domain: ${domain})`);
  230. // 固定的云函数ID
  231. const functionId = 'ZmYNPsoX9X';
  232. // 每页100条,共需要请求4次获取400条数据
  233. for (let page = 1; page <= 4; page++) {
  234. try {
  235. progress?.({ type: 'info', message: `类目 ${nodeId} 第 ${page}/4 页热销产品采集中...` });
  236. const requestBody = {
  237. path: "/api/CategoryProducts",
  238. method: "POST",
  239. query: { domain },
  240. body: {
  241. NodeId: nodeId,
  242. Page: page,
  243. Range: 400
  244. },
  245. functionId
  246. };
  247. const result = await relayClient.forwardSorftime({
  248. path: requestBody.path,
  249. method: requestBody.method,
  250. query: requestBody.query,
  251. body: requestBody.body
  252. });
  253. if (!relaySucceeded(result)) {
  254. throw new Error(result.message || `第${page}页数据采集失败`);
  255. }
  256. console.log(`[Sorftime Scheduler] 类目 ${nodeId} 第${page}页数据采集完成`);
  257. // 请求间隔,避免频率过高
  258. await this.#delay(2000);
  259. } catch (error) {
  260. console.error(`[Sorftime Scheduler] 类目 ${nodeId} 第${page}页采集失败:`, error.message);
  261. throw error;
  262. }
  263. }
  264. }
  265. /**
  266. * 采集类目关键词(每周一凌晨1点执行)
  267. * @private
  268. * @param {string} nodeId - 类目节点ID
  269. * @param {number} domain - 站点域名代码
  270. * @returns {Promise<void>}
  271. */
  272. async #collectCategoryKeywords(nodeId, domain, progress) {
  273. console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的关键词 (domain: ${domain})`);
  274. // 固定的云函数ID
  275. const functionId = 'KNN4L19aoi';
  276. try {
  277. const requestBody = {
  278. path: "/api/CategoryRequestKeyword",
  279. method: "POST",
  280. query: { domain },
  281. body: {
  282. Nodeid: nodeId,
  283. PageIndex: 1,
  284. PageSize: 100
  285. },
  286. functionId
  287. };
  288. const result = await relayClient.forwardSorftime({
  289. path: requestBody.path,
  290. method: requestBody.method,
  291. query: requestBody.query,
  292. body: requestBody.body
  293. });
  294. if (!relaySucceeded(result)) {
  295. throw new Error(result.message || '关键词采集失败');
  296. }
  297. console.log(`[Sorftime Scheduler] 类目 ${nodeId} 关键词采集完成`);
  298. } catch (error) {
  299. console.error(`[Sorftime Scheduler] 类目 ${nodeId} 关键词采集失败:`, error.message);
  300. throw error;
  301. }
  302. }
  303. /**
  304. * 通过 SellerId 分页采集单个店铺的产品数据
  305. * @private
  306. * @param {string} shopId - 店铺 Parse objectId
  307. * @param {string} sellerId - 亚马逊 SellerId
  308. * @param {number} domain - 站点代码(1=US)
  309. * @param {Function} progress - 可选进度回调
  310. * @returns {Promise<{processed: number, pages: number}>}
  311. */
  312. async #doCollectShopProducts(shopId, shopName, sellerId, domain, marketplaceId, progress) {
  313. const functionId = 'GSjAsvw9FK';
  314. let page = 1;
  315. let hasMore = true;
  316. let totalProcessed = 0;
  317. progress?.({ type: 'info', message: `SellerId: ${sellerId},开始分页采集产品数据...` });
  318. console.log(`[Sorftime Scheduler] 店铺 ${shopId} SellerId: ${sellerId} 开始采集`);
  319. while (hasMore) {
  320. progress?.({ type: 'info', message: `正在采集第 ${page} 页...` });
  321. console.log(`[Sorftime Scheduler] 店铺 ${shopId} 采集第 ${page} 页`);
  322. const result = await relayClient.forwardSorftime({
  323. path: '/api/ProductQuery',
  324. method: 'POST',
  325. query: { domain },
  326. body: {
  327. Query: 1,
  328. QueryType: 5,
  329. Pattern: sellerId,
  330. Page: page
  331. }
  332. });
  333. if (!relaySucceeded(result)) {
  334. throw new Error(`第 ${page} 页采集失败: ${result.message || result.Message || JSON.stringify(result)}`);
  335. }
  336. const payload = unwrapRelayData(result);
  337. const items = payload?.Items || payload?.Products || (Array.isArray(payload) ? payload : []);
  338. const count = Array.isArray(items) ? items.length : 0;
  339. let saved = 0;
  340. for (const item of items) {
  341. if (await upsertSorftimeProduct(item, shopId, shopName, item?.MarketplaceId || marketplaceId || '', domain)) saved++;
  342. }
  343. totalProcessed += count;
  344. progress?.({ type: 'info', message: `第 ${page} 页获取 ${count} 条产品,入库 ${saved} 条`, processed: totalProcessed });
  345. console.log(`[Sorftime Scheduler] 店铺 ${shopId} 第 ${page} 页获取 ${count} 条产品,入库 ${saved} 条,累计 ${totalProcessed} 条`);
  346. const pageCount = Number(payload?.PageCount || 0);
  347. if (count === 0 || (pageCount > 0 && page >= pageCount) || page >= 100) {
  348. hasMore = false;
  349. } else {
  350. page++;
  351. await this.#delay(2000);
  352. }
  353. }
  354. return { processed: totalProcessed, pages: page };
  355. }
  356. /**
  357. * 采集类目市场趋势(每月1号凌晨1点执行)
  358. * @private
  359. * @param {string} nodeId - 类目节点ID
  360. * @param {number} domain - 站点域名代码
  361. * @returns {Promise<void>}
  362. */
  363. async #collectCategoryMarketTrend(nodeId, domain, progress) {
  364. console.log(`[Sorftime Scheduler] 开始采集类目 ${nodeId} 的市场趋势 (domain: ${domain})`);
  365. const functionId = 'OCOKi93GjS';
  366. const trendIndexes = [0, 1, 2, 3, 4, 5];
  367. for (const TrendIndex of trendIndexes) {
  368. try {
  369. progress?.({ type: 'info', message: `类目 ${nodeId} 趋势类型 ${TrendIndex}/5 采集中...` });
  370. const requestBody = {
  371. path: "/api/CategoryTrend",
  372. method: "POST",
  373. query: { domain },
  374. body: {
  375. NodeId: nodeId,
  376. TrendIndex
  377. },
  378. functionId
  379. };
  380. const result = await relayClient.forwardSorftime({
  381. path: requestBody.path,
  382. method: requestBody.method,
  383. query: requestBody.query,
  384. body: requestBody.body
  385. });
  386. if (!relaySucceeded(result)) {
  387. throw new Error(result.message || `市场趋势类型 ${TrendIndex} 采集失败`);
  388. }
  389. console.log(`[Sorftime Scheduler] 类目 ${nodeId} 市场趋势类型 ${TrendIndex} 采集完成`);
  390. await this.#delay(1000);
  391. } catch (error) {
  392. console.error(`[Sorftime Scheduler] 类目 ${nodeId} 市场趋势类型 ${TrendIndex} 采集失败:`, error.message);
  393. throw error;
  394. }
  395. }
  396. }
  397. /**
  398. * 执行类目市场趋势采集任务(每月1号凌晨1点)
  399. * @private
  400. * @returns {Promise<any>} 执行结果对象
  401. */
  402. async #executeCategoryMarketTrendCollection(shopId, nodeIds, progress) {
  403. if (this.#isRunning) {
  404. console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
  405. return {
  406. success: false,
  407. message: '任务正在执行中',
  408. successCount: 0,
  409. errorCount: 0,
  410. duration: 0,
  411. errors: []
  412. };
  413. }
  414. this.#isRunning = true;
  415. const startTime = Date.now();
  416. const errors = [];
  417. let successCount = 0;
  418. let errorCount = 0;
  419. console.log('[Sorftime Scheduler] ========================================');
  420. console.log('[Sorftime Scheduler] 开始执行类目市场趋势采集任务');
  421. console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
  422. console.log('[Sorftime Scheduler] ========================================');
  423. try {
  424. const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
  425. if (categories.length === 0) {
  426. console.log('[Sorftime Scheduler] 没有找到叶子节点类目数据');
  427. return {
  428. success: true,
  429. message: '没有叶子节点类目需要处理',
  430. successCount: 0,
  431. errorCount: 0,
  432. duration: Date.now() - startTime,
  433. errors: []
  434. };
  435. }
  436. console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个叶子节点类目`);
  437. for (const category of categories) {
  438. try {
  439. const nodeId = category.get('nodeId');
  440. const categoryName = category.get('name');
  441. const domain = category.get('domain') || 1;
  442. console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
  443. console.log(`[Sorftime Scheduler] 开始处理类目市场趋势: ${categoryName} (${nodeId}, domain: ${domain})`);
  444. progress?.({ type: 'info', message: `开始采集类目【${categoryName}】市场趋势` });
  445. await this.#collectCategoryMarketTrend(nodeId, domain, progress);
  446. successCount++;
  447. progress?.({ type: 'success', message: `类目【${categoryName}】市场趋势采集完成`, processed: successCount });
  448. console.log(`[Sorftime Scheduler] 类目 ${categoryName} 市场趋势采集完成`);
  449. await this.#delay(2000);
  450. } catch (error) {
  451. errorCount++;
  452. const errorMsg = `类目 ${category.get('name')} 市场趋势采集失败: ${error.message}`;
  453. errors.push(errorMsg);
  454. progress?.({ type: 'error', message: errorMsg });
  455. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  456. }
  457. }
  458. const duration = Date.now() - startTime;
  459. await this.#logExecution({
  460. taskName: 'sorftime-category-trend-monthly',
  461. startTime: new Date(startTime),
  462. endTime: new Date(),
  463. duration,
  464. successCount,
  465. errorCount,
  466. errors,
  467. status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
  468. });
  469. this.#lastRunTime = new Date();
  470. const message = `类目市场趋势采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
  471. console.log(`\n[Sorftime Scheduler] ========================================`);
  472. console.log(`[Sorftime Scheduler] ${message}`);
  473. console.log(`[Sorftime Scheduler] ========================================\n`);
  474. return {
  475. success: errorCount === 0,
  476. message,
  477. successCount,
  478. errorCount,
  479. duration,
  480. errors
  481. };
  482. } catch (error) {
  483. const duration = Date.now() - startTime;
  484. const errorMsg = `类目市场趋势采集任务执行失败: ${error.message}`;
  485. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  486. return {
  487. success: false,
  488. message: errorMsg,
  489. successCount,
  490. errorCount: errorCount + 1,
  491. duration,
  492. errors: [...errors, errorMsg]
  493. };
  494. } finally {
  495. this.#isRunning = false;
  496. }
  497. }
  498. /**
  499. * 执行类目热销产品采集任务(每月1号)
  500. * @private
  501. * @returns {Promise<any>} 执行结果对象,包含成功数、失败数、耗时等信息
  502. */
  503. async #executeCategoryProductsCollection(shopId, nodeIds, progress) {
  504. if (this.#isRunning) {
  505. console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
  506. return {
  507. success: false,
  508. message: '任务正在执行中',
  509. successCount: 0,
  510. errorCount: 0,
  511. duration: 0,
  512. errors: []
  513. };
  514. }
  515. this.#isRunning = true;
  516. const startTime = Date.now();
  517. const errors = [];
  518. let successCount = 0;
  519. let errorCount = 0;
  520. console.log('[Sorftime Scheduler] ========================================');
  521. console.log('[Sorftime Scheduler] 开始执行类目热销产品采集任务');
  522. console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
  523. console.log('[Sorftime Scheduler] ========================================');
  524. try {
  525. // 获取所有类目
  526. const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
  527. if (categories.length === 0) {
  528. console.log('[Sorftime Scheduler] 没有找到类目数据');
  529. return {
  530. success: true,
  531. message: '没有类目需要处理',
  532. successCount: 0,
  533. errorCount: 0,
  534. duration: Date.now() - startTime,
  535. errors: []
  536. };
  537. }
  538. console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个类目`);
  539. // 循环处理每个类目
  540. for (const category of categories) {
  541. try {
  542. const nodeId = category.get('nodeId');
  543. const categoryName = category.get('name');
  544. const domain = category.get('domain') || 1; // 从类目中获取domain,默认为1(美国站)
  545. console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
  546. console.log(`[Sorftime Scheduler] 开始处理类目: ${categoryName} (${nodeId}, domain: ${domain})`);
  547. progress?.({ type: 'info', message: `开始采集类目【${categoryName}】热销产品` });
  548. await this.#collectCategoryProducts(nodeId, domain, progress);
  549. successCount++;
  550. progress?.({ type: 'success', message: `类目【${categoryName}】热销产品采集完成`, processed: successCount });
  551. console.log(`[Sorftime Scheduler] 类目 ${categoryName} 处理完成`);
  552. await this.#delay(3000);
  553. } catch (error) {
  554. errorCount++;
  555. const errorMsg = `类目 ${category.get('name')} 处理失败: ${error.message}`;
  556. errors.push(errorMsg);
  557. progress?.({ type: 'error', message: errorMsg });
  558. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  559. }
  560. }
  561. // 记录执行日志
  562. const duration = Date.now() - startTime;
  563. await this.#logExecution({
  564. taskName: 'sorftime-category-products-monthly',
  565. startTime: new Date(startTime),
  566. endTime: new Date(),
  567. duration,
  568. successCount,
  569. errorCount,
  570. errors,
  571. status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
  572. });
  573. this.#lastRunTime = new Date();
  574. const message = `类目热销产品采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
  575. console.log(`\n[Sorftime Scheduler] ========================================`);
  576. console.log(`[Sorftime Scheduler] ${message}`);
  577. console.log(`[Sorftime Scheduler] ========================================\n`);
  578. return {
  579. success: errorCount === 0,
  580. message,
  581. successCount,
  582. errorCount,
  583. duration,
  584. errors
  585. };
  586. } catch (error) {
  587. const duration = Date.now() - startTime;
  588. const errorMsg = `类目热销产品采集任务执行失败: ${error.message}`;
  589. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  590. return {
  591. success: false,
  592. message: errorMsg,
  593. successCount,
  594. errorCount: errorCount + 1,
  595. duration,
  596. errors: [...errors, errorMsg]
  597. };
  598. } finally {
  599. this.#isRunning = false;
  600. }
  601. }
  602. /**
  603. * 执行类目关键词采集任务(每周一)
  604. * @private
  605. * @returns {Promise<any>} 执行结果对象,包含成功数、失败数、耗时等信息
  606. */
  607. async #executeCategoryKeywordsCollection(shopId, nodeIds, progress) {
  608. if (this.#isRunning) {
  609. console.log('[Sorftime Scheduler] 任务正在执行中,跳过本次调度');
  610. return {
  611. success: false,
  612. message: '任务正在执行中',
  613. successCount: 0,
  614. errorCount: 0,
  615. duration: 0,
  616. errors: []
  617. };
  618. }
  619. this.#isRunning = true;
  620. const startTime = Date.now();
  621. const errors = [];
  622. let successCount = 0;
  623. let errorCount = 0;
  624. console.log('[Sorftime Scheduler] ========================================');
  625. console.log('[Sorftime Scheduler] 开始执行类目关键词采集任务');
  626. console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
  627. console.log('[Sorftime Scheduler] ========================================');
  628. try {
  629. // 获取所有叶子节点类目
  630. const categories = await this.#getLeafCategoriesByFilter(shopId, nodeIds);
  631. if (categories.length === 0) {
  632. console.log('[Sorftime Scheduler] 没有找到叶子节点类目数据');
  633. return {
  634. success: true,
  635. message: '没有叶子节点类目需要处理',
  636. successCount: 0,
  637. errorCount: 0,
  638. duration: Date.now() - startTime,
  639. errors: []
  640. };
  641. }
  642. console.log(`[Sorftime Scheduler] 发现 ${categories.length} 个叶子节点类目`);
  643. // 循环处理每个类目
  644. for (const category of categories) {
  645. try {
  646. const nodeId = category.get('nodeId');
  647. const categoryName = category.get('name');
  648. const domain = category.get('domain') || 1; // 从类目中获取domain,默认为1(美国站)
  649. console.log(`\n[Sorftime Scheduler] ----------------------------------------`);
  650. console.log(`[Sorftime Scheduler] 开始处理类目: ${categoryName} (${nodeId}, domain: ${domain})`);
  651. progress?.({ type: 'info', message: `开始采集类目【${categoryName}】关键词` });
  652. await this.#collectCategoryKeywords(nodeId, domain, progress);
  653. successCount++;
  654. progress?.({ type: 'success', message: `类目【${categoryName}】关键词采集完成`, processed: successCount });
  655. console.log(`[Sorftime Scheduler] 类目 ${categoryName} 关键词采集完成`);
  656. await this.#delay(2000);
  657. } catch (error) {
  658. errorCount++;
  659. const errorMsg = `类目 ${category.get('name')} 关键词采集失败: ${error.message}`;
  660. errors.push(errorMsg);
  661. progress?.({ type: 'error', message: errorMsg });
  662. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  663. }
  664. }
  665. // 记录执行日志
  666. const duration = Date.now() - startTime;
  667. await this.#logExecution({
  668. taskName: 'sorftime-category-keywords-weekly',
  669. startTime: new Date(startTime),
  670. endTime: new Date(),
  671. duration,
  672. successCount,
  673. errorCount,
  674. errors,
  675. status: errorCount === 0 ? 'success' : (successCount > 0 ? 'partial_success' : 'failed')
  676. });
  677. this.#lastRunTime = new Date();
  678. const message = `类目关键词采集完成: 成功 ${successCount}/${categories.length} 个类目,耗时 ${Math.round(duration / 1000)}秒`;
  679. console.log(`\n[Sorftime Scheduler] ========================================`);
  680. console.log(`[Sorftime Scheduler] ${message}`);
  681. console.log(`[Sorftime Scheduler] ========================================\n`);
  682. return {
  683. success: errorCount === 0,
  684. message,
  685. successCount,
  686. errorCount,
  687. duration,
  688. errors
  689. };
  690. } catch (error) {
  691. const duration = Date.now() - startTime;
  692. const errorMsg = `类目关键词采集任务执行失败: ${error.message}`;
  693. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  694. return {
  695. success: false,
  696. message: errorMsg,
  697. successCount,
  698. errorCount: errorCount + 1,
  699. duration,
  700. errors: [...errors, errorMsg]
  701. };
  702. } finally {
  703. this.#isRunning = false;
  704. }
  705. }
  706. /**
  707. * 执行云函数调用任务(每两小时)
  708. * 调用指定的云函数进行数据处理
  709. * @private
  710. * @returns {Promise<any>} 执行结果对象
  711. */
  712. async #executeCloudFunctionInvocation() {
  713. console.log('[Sorftime Scheduler] ========================================');
  714. console.log('[Sorftime Scheduler] 开始执行云函数调用任务');
  715. console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
  716. console.log('[Sorftime Scheduler] ========================================');
  717. const startTime = Date.now();
  718. try {
  719. const response = await fetch('http://localhost:3000/api/functions', {
  720. method: 'POST',
  721. headers: {
  722. 'Content-Type': 'application/json'
  723. },
  724. body: JSON.stringify({
  725. id: "Z4z6SB4o9e"
  726. })
  727. });
  728. const duration = Date.now() - startTime;
  729. if (!response.ok) {
  730. const errorMsg = `云函数调用失败,状态码: ${response.status}`;
  731. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  732. return {
  733. success: false,
  734. message: errorMsg,
  735. duration,
  736. errors: [errorMsg]
  737. };
  738. }
  739. const result = await response.json();
  740. const message = `云函数调用完成,耗时 ${Math.round(duration / 1000)}秒`;
  741. console.log(`\n[Sorftime Scheduler] ========================================`);
  742. console.log(`[Sorftime Scheduler] ${message}`);
  743. console.log(`[Sorftime Scheduler] ========================================\n`);
  744. return {
  745. success: true,
  746. message,
  747. duration,
  748. data: result,
  749. errors: []
  750. };
  751. } catch (error) {
  752. const duration = Date.now() - startTime;
  753. const errorMsg = `云函数调用任务执行失败: ${error.message}`;
  754. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  755. return {
  756. success: false,
  757. message: errorMsg,
  758. duration,
  759. errors: [errorMsg]
  760. };
  761. }
  762. }
  763. /**
  764. * 执行产品评论采集任务(每天凌晨1点)
  765. * 根据 Product 表中的 ASIN,采集对应的评论信息
  766. * @private
  767. * @returns {Promise<any>} 执行结果对象
  768. */
  769. async #executeProductReviewsCollection(progress) {
  770. console.log('[Sorftime Scheduler] ========================================');
  771. console.log('[Sorftime Scheduler] 开始执行产品评论采集任务');
  772. console.log('[Sorftime Scheduler] 执行时间:', new Date().toISOString());
  773. console.log('[Sorftime Scheduler] ========================================');
  774. const startTime = Date.now();
  775. let successCount = 0;
  776. let errorCount = 0;
  777. const errors = [];
  778. progress?.({ type: 'info', message: '开始执行产品评论采集任务' });
  779. try {
  780. const Parse = globalThis.Parse;
  781. // 查询所有 Product 记录
  782. const productQuery = new Parse.Query('Product');
  783. productQuery.limit(10000);
  784. const products = await productQuery.find({ useMasterKey: true });
  785. if (products.length === 0) {
  786. console.log('[Sorftime Scheduler] 没有找到 Product 数据');
  787. progress?.({ type: 'info', message: '没有 Product 数据需要处理' });
  788. return {
  789. success: true,
  790. message: '没有 Product 数据需要处理',
  791. successCount: 0,
  792. errorCount: 0,
  793. duration: Date.now() - startTime,
  794. errors: []
  795. };
  796. }
  797. console.log(`[Sorftime Scheduler] 发现 ${products.length} 个 Product 记录`);
  798. progress?.({ type: 'info', message: `共发现 ${products.length} 个产品需要检查评论` });
  799. // 遍历每个 Product
  800. for (const product of products) {
  801. try {
  802. const asin = product.get('asin');
  803. if (!asin) {
  804. console.log('[Sorftime Scheduler] 跳过无 ASIN 的 Product');
  805. continue;
  806. }
  807. progress?.({ type: 'info', message: `开始处理 ASIN: ${asin}` });
  808. // 获取 Product 关联的 Shop 信息
  809. const shopRelation = product.get('shop');
  810. if (!shopRelation) {
  811. console.log(`[Sorftime Scheduler] Product 无关联 Shop,ASIN: ${asin}`);
  812. continue;
  813. }
  814. const shop = await shopRelation.fetch({ useMasterKey: true });
  815. const shopDomain = shop.get('domain') || 1;
  816. // 查询 SorftimeReviews 中该 ASIN 最新的 updateAt 时间
  817. let queryStartDt = '2025-01-01';
  818. const reviewQuery = new Parse.Query('SorftimeReviews');
  819. reviewQuery.equalTo('asin', asin);
  820. reviewQuery.descending('updatedAt');
  821. reviewQuery.limit(1);
  822. let latestReview = null;
  823. try {
  824. latestReview = await reviewQuery.first({ useMasterKey: true });
  825. } catch (error) {
  826. const message = String(error?.message || error || '');
  827. if (!message.includes('does not exist') && !message.includes('non-existent class')) throw error;
  828. }
  829. if (latestReview) {
  830. const updateAt = latestReview.get('updatedAt');
  831. if (updateAt) {
  832. // 格式化为 YYYY-MM-DD
  833. const date = new Date(updateAt);
  834. queryStartDt = date.toISOString().split('T')[0];
  835. }
  836. }
  837. console.log(`[Sorftime Scheduler] 开始采集评论,ASIN: ${asin}, 查询起始日期: ${queryStartDt}`);
  838. // 调用 Sorftime 接口获取评论信息
  839. let result;
  840. try {
  841. result = await relayClient.forwardSorftime({
  842. path: '/api/ProductReviewsQuery',
  843. method: 'POST',
  844. query: {
  845. domain: shopDomain,
  846. shopId: shop.id
  847. },
  848. body: {
  849. ASIN: asin,
  850. PageIndex: 1,
  851. OnlyPurchase: 1,
  852. Star: '1,2,3,4,5',
  853. QueryStartDt: queryStartDt
  854. }
  855. });
  856. } catch (error) {
  857. const errorMsg = `请求评论数据失败,ASIN: ${asin}, ${error.message || error}`;
  858. errors.push(errorMsg);
  859. errorCount++;
  860. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  861. continue;
  862. }
  863. const payload = unwrapRelayData(result);
  864. const reviews = payload?.Reviews || payload?.Items || payload?.List || (Array.isArray(payload) ? payload : []);
  865. if (relaySucceeded(result)) {
  866. for (const review of reviews) await upsertSorftimeReview(review, asin, shop.id);
  867. successCount++;
  868. const statusText = reviews.length ? `入库 ${reviews.length} 条` : '无新增评论';
  869. console.log(`[Sorftime Scheduler] 评论数据采集成功,ASIN: ${asin},${statusText}`);
  870. progress?.({ type: 'success', message: `评论数据采集成功,ASIN: ${asin},${statusText}`, processed: successCount });
  871. } else {
  872. const errorMsg = `评论数据获取失败,ASIN: ${asin}`;
  873. errors.push(errorMsg);
  874. errorCount++;
  875. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  876. progress?.({ type: 'error', message: errorMsg });
  877. }
  878. // 延迟 1 秒,避免请求过于频繁
  879. await this.#delay(500);
  880. } catch (error) {
  881. errorCount++;
  882. const errorMsg = `处理 Product 失败: ${error.message}`;
  883. errors.push(errorMsg);
  884. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  885. progress?.({ type: 'error', message: errorMsg });
  886. }
  887. }
  888. const duration = Date.now() - startTime;
  889. const message = `产品评论采集完成: 成功 ${successCount}/${products.length} 个产品,耗时 ${Math.round(duration / 1000)}秒`;
  890. console.log(`\n[Sorftime Scheduler] ========================================`);
  891. console.log(`[Sorftime Scheduler] ${message}`);
  892. console.log(`[Sorftime Scheduler] ========================================\n`);
  893. progress?.({ type: errorCount === 0 ? 'success' : 'warning', message });
  894. return {
  895. success: errorCount === 0,
  896. message,
  897. successCount,
  898. errorCount,
  899. duration,
  900. errors
  901. };
  902. } catch (error) {
  903. const duration = Date.now() - startTime;
  904. const errorMsg = `产品评论采集任务执行失败: ${error.message}`;
  905. console.error(`[Sorftime Scheduler] ${errorMsg}`);
  906. progress?.({ type: 'error', message: errorMsg });
  907. return {
  908. success: false,
  909. message: errorMsg,
  910. successCount: 0,
  911. errorCount: 1,
  912. duration,
  913. errors: [errorMsg]
  914. };
  915. }
  916. }
  917. // ========== 公共方法(后声明,因为依赖前面的私有方法) ==========
  918. /**
  919. * 构造函数
  920. * 初始化 Sorftime API 数据采集调度器
  921. */
  922. constructor() {
  923. console.log('[Sorftime Scheduler] 初始化 Sorftime API 数据采集调度器');
  924. }
  925. /**
  926. * 启动定时任务
  927. * @returns {Promise<void>}
  928. */
  929. async start() {
  930. if (this.#categoryProductsCronTask || this.#keywordsCronTask || this.#marketTrendCronTask || this.#shopProductsCronTask || this.#productDetailCronTask || this.#cloudFunctionCronTask || this.#reviewsCronTask) {
  931. console.log('[Sorftime Scheduler] 定时任务已在运行中');
  932. return;
  933. }
  934. // 每月1号凌晨执行类目热销产品采集
  935. this.#categoryProductsCronTask = nodeCron.schedule('0 0 1 * *', async () => {
  936. await this.#executeCategoryProductsCollection(undefined, undefined, undefined);
  937. }, {
  938. timezone: "Asia/Shanghai"
  939. });
  940. console.log('[Sorftime Scheduler] 类目热销产品定时任务已启动,将在每月1号凌晨执行');
  941. // 每月1号凌晨1点执行类目市场趋势采集
  942. this.#marketTrendCronTask = nodeCron.schedule('0 1 1 * *', async () => {
  943. await this.#executeCategoryMarketTrendCollection(undefined, undefined, undefined);
  944. }, {
  945. timezone: "Asia/Shanghai"
  946. });
  947. console.log('[Sorftime Scheduler] 类目市场趋势定时任务已启动,将在每月1号凌晨1点执行');
  948. // 每周一凌晨1点执行类目关键词采集
  949. this.#keywordsCronTask = nodeCron.schedule('0 1 * * 1', async () => {
  950. await this.#executeCategoryKeywordsCollection(undefined, undefined, undefined);
  951. }, {
  952. timezone: "Asia/Shanghai"
  953. });
  954. console.log('[Sorftime Scheduler] 类目关键词定时任务已启动,将在每周一凌晨1点执行');
  955. if (process.env.ENABLE_SORFTIME_REVIEW_SCHEDULES === 'true') {
  956. this.#reviewsCronTask = nodeCron.schedule('0 1 * * *', async () => {
  957. console.log('[Sorftime Scheduler] ===== 凌晨1点:开始产品评论采集 =====');
  958. try {
  959. await this.#executeProductReviewsCollection(undefined);
  960. } catch (e) {
  961. console.error('[Sorftime Scheduler] 产品评论采集失败:', e.message);
  962. }
  963. console.log('[Sorftime Scheduler] ===== 产品评论采集完成 =====');
  964. }, { timezone: 'Asia/Shanghai' });
  965. console.log('[Sorftime Scheduler] 产品评论采集定时任务已启动,将在每天凌晨1点执行');
  966. } else {
  967. console.log('[Sorftime Scheduler] 产品评论采集定时任务已禁用');
  968. }
  969. // 每天凌晨2点采集店铺产品数据
  970. this.#shopProductsCronTask = nodeCron.schedule('0 2 * * *', async () => {
  971. console.log('[Sorftime Scheduler] ===== 凌晨2点:开始店铺产品采集 =====');
  972. try {
  973. await this.collectShopProducts(undefined, undefined);
  974. } catch (e) {
  975. console.error('[Sorftime Scheduler] 店铺产品采集失败:', e.message);
  976. }
  977. console.log('[Sorftime Scheduler] ===== 店铺产品采集完成 =====');
  978. }, { timezone: 'Asia/Shanghai' });
  979. console.log('[Sorftime Scheduler] 店铺产品采集定时任务已启动,将在每天凌晨2点执行');
  980. if (process.env.ENABLE_LEGACY_CLOUD_FUNCTION_SCHEDULES === 'true') {
  981. this.#cloudFunctionCronTask = nodeCron.schedule('0 */2 * * *', async () => {
  982. try {
  983. await this.#executeCloudFunctionInvocation();
  984. } catch (e) {
  985. console.error('[Sorftime Scheduler] 遗留云函数调用失败:', e.message);
  986. }
  987. }, { timezone: 'Asia/Shanghai' });
  988. } else {
  989. console.log('[Sorftime Scheduler] 遗留动态云函数定时任务已禁用');
  990. }
  991. }
  992. /**
  993. * 停止定时任务
  994. * @returns {void}
  995. */
  996. stop() {
  997. if (this.#categoryProductsCronTask) {
  998. this.#categoryProductsCronTask.stop();
  999. this.#categoryProductsCronTask = null;
  1000. console.log('[Sorftime Scheduler] 类目热销产品定时任务已停止');
  1001. }
  1002. if (this.#marketTrendCronTask) {
  1003. this.#marketTrendCronTask.stop();
  1004. this.#marketTrendCronTask = null;
  1005. console.log('[Sorftime Scheduler] 类目市场趋势定时任务已停止');
  1006. }
  1007. if (this.#keywordsCronTask) {
  1008. this.#keywordsCronTask.stop();
  1009. this.#keywordsCronTask = null;
  1010. console.log('[Sorftime Scheduler] 类目关键词定时任务已停止');
  1011. }
  1012. if (this.#shopProductsCronTask) {
  1013. this.#shopProductsCronTask.stop();
  1014. this.#shopProductsCronTask = null;
  1015. console.log('[Sorftime Scheduler] 店铺产品采集定时任务已停止');
  1016. }
  1017. if (this.#productDetailCronTask) {
  1018. this.#productDetailCronTask.stop();
  1019. this.#productDetailCronTask = null;
  1020. console.log('[Sorftime Scheduler] 产品详情同步定时任务已停止');
  1021. }
  1022. if (this.#cloudFunctionCronTask) {
  1023. this.#cloudFunctionCronTask.stop();
  1024. this.#cloudFunctionCronTask = null;
  1025. console.log('[Sorftime Scheduler] 云函数调用定时任务已停止');
  1026. }
  1027. if (this.#reviewsCronTask) {
  1028. this.#reviewsCronTask.stop();
  1029. this.#reviewsCronTask = null;
  1030. console.log('[Sorftime Scheduler] 产品评论采集定时任务已停止');
  1031. }
  1032. }
  1033. /**
  1034. * 获取调度器状态
  1035. * @returns {object} 调度器状态对象
  1036. */
  1037. getStatus() {
  1038. return {
  1039. isRunning: this.#isRunning,
  1040. lastRunTime: this.#lastRunTime,
  1041. cronEnabled: this.#cronEnabled,
  1042. tasks: {
  1043. categoryProducts: {
  1044. enabled: !!this.#categoryProductsCronTask,
  1045. schedule: '每月1号凌晨'
  1046. },
  1047. marketTrend: {
  1048. enabled: !!this.#marketTrendCronTask,
  1049. schedule: '每月1号凌晨1点'
  1050. },
  1051. keywords: {
  1052. enabled: !!this.#keywordsCronTask,
  1053. schedule: '每周一凌晨1点'
  1054. },
  1055. shopProducts: {
  1056. enabled: !!this.#shopProductsCronTask,
  1057. schedule: '每天凌晨2点'
  1058. },
  1059. productDetails: {
  1060. enabled: !!this.#productDetailCronTask,
  1061. schedule: '每小时'
  1062. },
  1063. cloudFunction: {
  1064. enabled: !!this.#cloudFunctionCronTask,
  1065. schedule: '每两小时'
  1066. },
  1067. productReviews: {
  1068. enabled: !!this.#reviewsCronTask,
  1069. schedule: '每天凌晨1点'
  1070. }
  1071. }
  1072. };
  1073. }
  1074. /**
  1075. * 手动触发类目热销产品采集
  1076. * @returns {Promise<any>} 执行结果对象
  1077. */
  1078. async triggerCategoryProducts(shopId, nodeIds, progress) {
  1079. console.log('[Sorftime Scheduler] 手动触发类目热销产品采集任务');
  1080. return await this.#executeCategoryProductsCollection(shopId, nodeIds, progress);
  1081. }
  1082. /**
  1083. * 手动触发类目关键词采集
  1084. * @returns {Promise<any>} 执行结果对象
  1085. */
  1086. async triggerCategoryKeywords(shopId, nodeIds, progress) {
  1087. console.log('[Sorftime Scheduler] 手动触发类目关键词采集任务');
  1088. return await this.#executeCategoryKeywordsCollection(shopId, nodeIds, progress);
  1089. }
  1090. /**
  1091. * 手动触发类目市场趋势采集
  1092. * @returns {Promise<any>} 执行结果对象
  1093. */
  1094. async triggerCategoryTrends(shopId, nodeIds, progress) {
  1095. console.log('[Sorftime Scheduler] 手动触发类目市场趋势采集任务');
  1096. return await this.#executeCategoryMarketTrendCollection(shopId, nodeIds, progress);
  1097. }
  1098. /**
  1099. * 手动触发产品评论采集
  1100. * @returns {Promise<any>} 执行结果对象
  1101. */
  1102. async triggerProductReviews(progress) {
  1103. console.log('[Sorftime Scheduler] 手动触发产品评论采集任务');
  1104. return await this.#executeProductReviewsCollection(progress);
  1105. }
  1106. /**
  1107. * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势)
  1108. * 前端选中某个类目但数据库无数据时调用
  1109. * @param {string} nodeId - 类目节点ID
  1110. * @param {number} domain - 站点域名代码(默认1=美国)
  1111. * @returns {Promise<any>} 执行结果对象
  1112. */
  1113. async triggerSingleCategory(nodeId, domain = 1, progress) {
  1114. const startTime = Date.now();
  1115. const errors = [];
  1116. let tasks = { products: false, keywords: false, trends: false };
  1117. console.log(`[Sorftime Scheduler] 按需触发单类目采集: nodeId=${nodeId}, domain=${domain}`);
  1118. try {
  1119. // 1. 热销产品
  1120. try {
  1121. progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 热销产品` });
  1122. await this.#collectCategoryProducts(nodeId, domain, progress);
  1123. tasks.products = true;
  1124. progress?.({ type: 'success', message: `类目 ${nodeId} 热销产品采集完成` });
  1125. console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 热销产品采集完成`);
  1126. } catch (error) {
  1127. const errMsg = error?.message || String(error);
  1128. progress?.({ type: 'error', message: `热销产品采集失败: ${errMsg}` });
  1129. errors.push(`热销产品采集失败: ${errMsg}`);
  1130. console.error(`[Sorftime Scheduler] 热销产品采集失败:`, error);
  1131. }
  1132. // 延迟 2 秒
  1133. await new Promise(resolve => setTimeout(resolve, 500));
  1134. // 2. 关键词
  1135. try {
  1136. progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 关键词` });
  1137. await this.#collectCategoryKeywords(nodeId, domain, progress);
  1138. tasks.keywords = true;
  1139. progress?.({ type: 'success', message: `类目 ${nodeId} 关键词采集完成` });
  1140. console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 关键词采集完成`);
  1141. } catch (error) {
  1142. const errMsg = error?.message || String(error);
  1143. progress?.({ type: 'error', message: `关键词采集失败: ${errMsg}` });
  1144. errors.push(`关键词采集失败: ${errMsg}`);
  1145. console.error(`[Sorftime Scheduler] 关键词采集失败:`, error);
  1146. }
  1147. // 延迟 2 秒
  1148. await new Promise(resolve => setTimeout(resolve, 500));
  1149. // 3. 市场趋势
  1150. try {
  1151. progress?.({ type: 'info', message: `开始采集类目 ${nodeId} 市场趋势` });
  1152. await this.#collectCategoryMarketTrend(nodeId, domain, progress);
  1153. tasks.trends = true;
  1154. progress?.({ type: 'success', message: `类目 ${nodeId} 市场趋势采集完成` });
  1155. console.log(`[Sorftime Scheduler] 单类目 ${nodeId} 市场趋势采集完成`);
  1156. } catch (error) {
  1157. const errMsg = error?.message || String(error);
  1158. progress?.({ type: 'error', message: `市场趋势采集失败: ${errMsg}` });
  1159. errors.push(`市场趋势采集失败: ${errMsg}`);
  1160. console.error(`[Sorftime Scheduler] 市场趋势采集失败:`, error);
  1161. }
  1162. const duration = Date.now() - startTime;
  1163. const allSuccess = errors.length === 0;
  1164. const successCount = Object.values(tasks).filter(Boolean).length;
  1165. // 异步记录执行日志,不阻塞返回
  1166. if (this.#logExecution) {
  1167. this.#logExecution({
  1168. taskName: `trigger-single-category-${nodeId}`,
  1169. startTime: new Date(startTime),
  1170. endTime: new Date(),
  1171. duration,
  1172. successCount,
  1173. errorCount: errors.length,
  1174. errors,
  1175. status: allSuccess ? 'success' : 'partial_success'
  1176. }).catch(err => console.error('[Sorftime Scheduler] Log execution error:', err));
  1177. }
  1178. const result = {
  1179. success: allSuccess,
  1180. message: allSuccess
  1181. ? `类目 ${nodeId} 全量采集完成,耗时 ${Math.round(duration / 1000)}秒`
  1182. : `类目 ${nodeId} 部分采集失败(成功 ${successCount}/3 个任务)`,
  1183. nodeId,
  1184. domain,
  1185. tasks,
  1186. duration,
  1187. errors: errors.length > 0 ? errors : undefined
  1188. };
  1189. console.log(`[Sorftime Scheduler] 采集结果:`, result);
  1190. return result;
  1191. } catch (err) {
  1192. const duration = Date.now() - startTime;
  1193. const errMsg = err?.message || String(err);
  1194. console.error(`[Sorftime Scheduler] 采集异常:`, err);
  1195. return {
  1196. success: false,
  1197. message: `类目 ${nodeId} 采集异常: ${errMsg}`,
  1198. nodeId,
  1199. domain,
  1200. tasks,
  1201. duration,
  1202. errors: [errMsg]
  1203. };
  1204. }
  1205. }
  1206. /**
  1207. * 采集所有活跃 Amazon 店铺(或指定店铺)的产品数据
  1208. * 通过 SellerId (QueryType=5) 分页请求 /api/ProductQuery
  1209. * @param {string} shopId - 可选,不传则采集全部活跃 Amazon 店铺
  1210. * @param {Function} progress - 可选进度回调,适合 SSE 实时推送
  1211. * @returns {Promise<any>} 执行结果
  1212. */
  1213. async collectShopProducts(shopId, progress) {
  1214. const shops = await this.#getActiveAmazonShops(shopId);
  1215. if (shops.length === 0) {
  1216. const msg = shopId ? `未找到店铺 ${shopId} 或店铺未激活` : '没有找到活跃的 Amazon 店铺';
  1217. progress?.({ type: 'warn', message: msg });
  1218. console.warn(`[Sorftime Scheduler] ${msg}`);
  1219. return { success: true, message: msg, processed: 0, errorCount: 0, errors: [] };
  1220. }
  1221. let totalProcessed = 0;
  1222. let errorCount = 0;
  1223. const errors = [];
  1224. for (const shop of shops) {
  1225. const sid = shop.id;
  1226. const shopName = shop.get('name');
  1227. const config = shop.get('config');
  1228. const domain = Number(shop.get('domain') || 1);
  1229. const marketplaceId = shop.get('marketplaceId') || '';
  1230. if (!config || !config.SpApiConfig || !config.SpApiConfig.sellerID) {
  1231. const msg = `店铺【${shopName}】(${sid}) 缺少 SpApiConfig.sellerID 配置,跳过`;
  1232. progress?.({ type: 'error', message: msg });
  1233. console.error(`[Sorftime Scheduler] ${msg}`);
  1234. errors.push(msg);
  1235. errorCount++;
  1236. continue;
  1237. }
  1238. if (config.SpApiConfig.listingEnabled === false) {
  1239. progress?.({ type: 'warn', message: `店铺【${shopName}】未配置区域 Seller ID,跳过 Sorftime 店铺产品采集` });
  1240. continue;
  1241. }
  1242. const sellerId = config.SpApiConfig.sellerID;
  1243. progress?.({ type: 'info', message: `开始采集店铺【${shopName}】的产品数据 (SellerId: ${sellerId})` });
  1244. console.log(`[Sorftime Scheduler] 开始采集店铺 ${shopName} (${sid}) 产品数据`);
  1245. try {
  1246. const result = await this.#doCollectShopProducts(sid, shopName, sellerId, domain, marketplaceId, progress);
  1247. progress?.({ type: 'success', message: `店铺【${shopName}】产品采集完成: 共 ${result.pages} 页 ${result.processed} 条`, processed: result.processed });
  1248. console.log(`[Sorftime Scheduler] 店铺 ${shopName} 产品采集完成: ${result.processed} 条`);
  1249. totalProcessed += result.processed;
  1250. await this.#delay(2000);
  1251. } catch (error) {
  1252. const msg = `店铺【${shopName}】产品采集失败: ${error.message}`;
  1253. progress?.({ type: 'error', message: msg });
  1254. console.error(`[Sorftime Scheduler] ${msg}`);
  1255. errors.push(msg);
  1256. errorCount++;
  1257. }
  1258. }
  1259. const message = `店铺产品采集完成: 成功 ${shops.length - errorCount}/${shops.length} 个店铺,共 ${totalProcessed} 条记录`;
  1260. return { success: errorCount === 0, message, processed: totalProcessed, errorCount, errors };
  1261. }
  1262. }
  1263. // 导出单例
  1264. export const sorftimeScheduler = new SorftimeScheduler();