sorftime-api-schedule.ts 46 KB

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