routes-schedule.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402
  1. import express from "npm:express";
  2. import { spApiScheduler } from './sp-api-schedule.ts';
  3. import { newSpApiScheduler } from './new-sp-api-schedule.ts';
  4. import { sorftimeScheduler } from './sorftime-api-schedule.ts';
  5. /**
  6. * 定时器清洗数据
  7. */
  8. (async () => {
  9. try {
  10. await spApiScheduler.start();
  11. console.log('[API Routes] Amazon SP-API定时任务已启动');
  12. } catch (error) {
  13. console.error('[API Routes] Amazon SP-API定时任务启动失败:', error);
  14. }
  15. try {
  16. await sorftimeScheduler.start();
  17. console.log('[API Routes] Sorftime定时任务已启动');
  18. } catch (error) {
  19. console.error('[API Routes] Sorftime定时任务启动失败:', error);
  20. }
  21. })();
  22. const router = express.Router();
  23. console.log('加载schedule路由 /api/schedule/');
  24. router.use(express.json({
  25. charset: 'utf-8',
  26. // 额外配置:确保解析URL编码的请求体也支持中文
  27. type: 'application/json'
  28. }));
  29. // 定时任务管理路由
  30. router.get('/sp-api-schedule/status', (req, res) => {
  31. res.json({
  32. success: true,
  33. data: {
  34. spApi: spApiScheduler.getStatus()
  35. }
  36. });
  37. });
  38. router.get('/sp-api-schedule/start', async (req, res) => {
  39. console.log('[API Routes] 启动定时任务');
  40. try {
  41. const result = await spApiScheduler.start()
  42. res.json({
  43. success: true,
  44. message: '定时器执行陈功',
  45. data: result
  46. });
  47. } catch (error) {
  48. res.status(500).json({
  49. success: false,
  50. message: error.message
  51. });
  52. }
  53. });
  54. router.post('/sp-api-schedule/trigger', async (req, res) => {
  55. try {
  56. const shopId = req.body.shopId;
  57. const result = await spApiScheduler.triggerManually(shopId);
  58. res.json({
  59. success: result.success,
  60. message: result.message,
  61. data: result
  62. });
  63. } catch (error) {
  64. res.status(500).json({
  65. success: false,
  66. message: error.message
  67. });
  68. }
  69. });
  70. // ===== 以下为新版独立采集接口(SSE 实时进度推送)=====
  71. // 前端通过 EventSource 或 fetch + ReadableStream 消费 SSE 流
  72. // 事件格式:
  73. // event: progress data: { type, message, processed? }
  74. // event: done data: { success, message, processed, errorCount, errors }
  75. // event: error data: { success: false, message }
  76. // 心跳注释行(': heartbeat')每 15 秒发送一次,防止连接超时
  77. /**
  78. * 创建 SSE 响应并返回辅助函数
  79. */
  80. function createSseResponse(res) {
  81. res.setHeader('Content-Type', 'text/event-stream; charset=utf-8');
  82. res.setHeader('Cache-Control', 'no-cache');
  83. res.setHeader('Connection', 'keep-alive');
  84. res.setHeader('X-Accel-Buffering', 'no');
  85. res.setHeader('Access-Control-Allow-Origin', '*');
  86. res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
  87. res.setHeader('Access-Control-Allow-Headers', 'Content-Type');
  88. res.flushHeaders();
  89. const send = (event, data) => {
  90. if (!res.writableEnded) {
  91. res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
  92. }
  93. };
  94. const heartbeatTimer = setInterval(() => {
  95. if (!res.writableEnded) res.write(': heartbeat\n\n');
  96. }, 15000);
  97. const close = () => {
  98. clearInterval(heartbeatTimer);
  99. if (!res.writableEnded) res.end();
  100. };
  101. return { send, close };
  102. }
  103. /**
  104. * POST /api/schedule/sp-api-schedule/collect-listing
  105. * 采集 Listing 数据,通过 SSE 流式返回进度
  106. * Body: { shopId: string }
  107. */
  108. router.post('/sp-api-schedule/collect-listing', async (req, res) => {
  109. const { shopId } = req.body || {};
  110. const { send, close } = createSseResponse(res);
  111. req.on('close', () => close());
  112. try {
  113. send('progress', { type: 'info', message: `开始采集 Listing 数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  114. const result = await newSpApiScheduler.collectListings(
  115. shopId || undefined,
  116. (event) => send('progress', event)
  117. );
  118. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  119. } catch (error) {
  120. send('error', { success: false, message: error.message || '采集 Listing 数据时发生未知错误' });
  121. } finally {
  122. close();
  123. }
  124. });
  125. /**
  126. * POST /api/schedule/sp-api-schedule/collect-order
  127. * 采集订单数据,通过 SSE 流式返回进度
  128. * Body: { shopId: string }
  129. */
  130. router.post('/sp-api-schedule/collect-order', async (req, res) => {
  131. const { shopId } = req.body || {};
  132. const { send, close } = createSseResponse(res);
  133. req.on('close', () => close());
  134. try {
  135. send('progress', { type: 'info', message: `开始采集订单数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  136. const result = await newSpApiScheduler.collectOrders(
  137. shopId || undefined,
  138. (event) => send('progress', event)
  139. );
  140. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  141. } catch (error) {
  142. send('error', { success: false, message: error.message || '采集订单数据时发生未知错误' });
  143. } finally {
  144. close();
  145. }
  146. });
  147. /**
  148. * POST /api/schedule/sp-api-schedule/collect-report
  149. * 采集报表数据(创建报表任务),通过 SSE 流式返回进度
  150. * Body: { shopId: string }
  151. */
  152. router.post('/sp-api-schedule/collect-report', async (req, res) => {
  153. const { shopId } = req.body || {};
  154. const { send, close } = createSseResponse(res);
  155. req.on('close', () => close());
  156. try {
  157. send('progress', { type: 'info', message: `开始采集报表数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  158. const result = await newSpApiScheduler.collectReports(
  159. shopId || undefined,
  160. (event) => send('progress', event)
  161. );
  162. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  163. } catch (error) {
  164. send('error', { success: false, message: error.message || '采集报表数据时发生未知错误' });
  165. } finally {
  166. close();
  167. }
  168. });
  169. /**
  170. * POST /api/schedule/sp-api-schedule/parse-report
  171. * 检查未完成报表状态并解析已完成报表,通过 SSE 流式返回进度
  172. * Body: { shopId: string }
  173. */
  174. router.post('/sp-api-schedule/parse-report', async (req, res) => {
  175. const { shopId } = req.body || {};
  176. const { send, close } = createSseResponse(res);
  177. req.on('close', () => close());
  178. try {
  179. send('progress', { type: 'info', message: `开始解析报表${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  180. const result = await newSpApiScheduler.parseReports(
  181. shopId || undefined,
  182. (event) => send('progress', event)
  183. );
  184. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  185. } catch (error) {
  186. send('error', { success: false, message: error.message || '解析报表时发生未知错误' });
  187. } finally {
  188. close();
  189. }
  190. });
  191. // Sorftime 定时任务管理路由
  192. router.get('/sorftime-schedule/status', (req, res) => {
  193. res.json({
  194. success: true,
  195. data: {
  196. sorftime: sorftimeScheduler.getStatus()
  197. }
  198. });
  199. });
  200. /**
  201. * POST /api/schedule/sorftime-schedule/trigger-products
  202. * 手动触发全量类目热销产品采集(每类目分4页请求),SSE 实时返回进度
  203. * 对应定时任务:每月1日 00:00(Asia/Shanghai)
  204. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  205. * SSE done: { success, message, successCount, errorCount, errors }
  206. */
  207. router.post('/sorftime-schedule/trigger-products', async (req, res) => {
  208. const { shopId, nodeIds } = req.body || {};
  209. const { send, close } = createSseResponse(res);
  210. req.on('close', () => close());
  211. try {
  212. send('progress', { type: 'info', message: '开始采集类目热销产品...' });
  213. const result = await sorftimeScheduler.triggerCategoryProducts(
  214. shopId || undefined,
  215. nodeIds || undefined,
  216. (event) => send('progress', event)
  217. );
  218. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  219. } catch (error) {
  220. send('error', { success: false, message: error.message || '采集类目热销产品时发生未知错误' });
  221. } finally {
  222. close();
  223. }
  224. });
  225. /**
  226. * POST /api/schedule/sorftime-schedule/trigger-keywords
  227. * 手动触发全量类目关键词采集,SSE 实时返回进度
  228. * 对应定时任务:每周一 01:00(Asia/Shanghai)
  229. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  230. * SSE done: { success, message, successCount, errorCount, errors }
  231. */
  232. router.post('/sorftime-schedule/trigger-keywords', async (req, res) => {
  233. const { shopId, nodeIds } = req.body || {};
  234. const { send, close } = createSseResponse(res);
  235. req.on('close', () => close());
  236. try {
  237. send('progress', { type: 'info', message: '开始采集类目关键词...' });
  238. const result = await sorftimeScheduler.triggerCategoryKeywords(
  239. shopId || undefined,
  240. nodeIds || undefined,
  241. (event) => send('progress', event)
  242. );
  243. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  244. } catch (error) {
  245. send('error', { success: false, message: error.message || '采集类目关键词时发生未知错误' });
  246. } finally {
  247. close();
  248. }
  249. });
  250. /**
  251. * POST /api/schedule/sorftime-schedule/trigger-category-trends
  252. * 手动触发全量类目市场趋势采集(每类目6种 TrendIndex),SSE 实时返回进度
  253. * 对应定时任务:每月1日 01:00(Asia/Shanghai)
  254. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  255. * SSE done: { success, message, successCount, errorCount, errors }
  256. */
  257. router.post('/sorftime-schedule/trigger-category-trends', async (req, res) => {
  258. const { shopId, nodeIds } = req.body || {};
  259. const { send, close } = createSseResponse(res);
  260. req.on('close', () => close());
  261. try {
  262. send('progress', { type: 'info', message: '开始采集类目市场趋势...' });
  263. const result = await sorftimeScheduler.triggerCategoryTrends(
  264. shopId || undefined,
  265. nodeIds || undefined,
  266. (event) => send('progress', event)
  267. );
  268. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  269. } catch (error) {
  270. send('error', { success: false, message: error.message || '采集类目市场趋势时发生未知错误' });
  271. } finally {
  272. close();
  273. }
  274. });
  275. /**
  276. * POST /api/schedule/sorftime-schedule/trigger-single-category
  277. * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势),SSE 实时返回进度
  278. * Body: { nodeId: string, domain?: number } domain 默认 1(美国站)
  279. * SSE done: { success, message, tasks: { products, keywords, trends }, duration, errors }
  280. */
  281. router.post('/sorftime-schedule/trigger-single-category', async (req, res) => {
  282. const { nodeId, domain } = req.body || {};
  283. if (!nodeId) {
  284. return res.status(400).json({ success: false, message: '缺少 nodeId 参数' });
  285. }
  286. const { send, close } = createSseResponse(res);
  287. req.on('close', () => close());
  288. try {
  289. send('progress', { type: 'info', message: `开始采集类目 ${nodeId} 全量数据(热销产品 + 关键词 + 市场趋势)` });
  290. const result = await sorftimeScheduler.triggerSingleCategory(
  291. nodeId,
  292. domain || 1,
  293. (event) => send('progress', event)
  294. );
  295. send('done', { success: result.success, message: result.message, tasks: result.tasks, duration: result.duration, errors: result.errors });
  296. } catch (error) {
  297. send('error', { success: false, message: error.message || '采集单类目数据时发生未知错误' });
  298. } finally {
  299. close();
  300. }
  301. });
  302. /**
  303. * POST /api/schedule/sorftime-schedule/collect-reviews
  304. * 手动触发产品评论采集(遍历 Product 表 ASIN 请求 ProductReviewsQuery),SSE 实时返回进度
  305. * Body: 无
  306. * SSE done: { success, message, successCount, errorCount, errors }
  307. */
  308. router.post('/sorftime-schedule/collect-reviews', async (req, res) => {
  309. const { send, close } = createSseResponse(res);
  310. req.on('close', () => close());
  311. try {
  312. send('progress', { type: 'info', message: '开始采集产品评论数据...' });
  313. const result = await sorftimeScheduler.triggerProductReviews(
  314. (event) => send('progress', event)
  315. );
  316. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  317. } catch (error) {
  318. send('error', { success: false, message: error.message || '采集产品评论时发生未知错误' });
  319. } finally {
  320. close();
  321. }
  322. });
  323. /**
  324. * POST /api/schedule/sorftime-schedule/collect-shop-products
  325. * 采集店铺竞品数据(通过 SellerId QueryType=5 分页请求),SSE 实时返回进度
  326. * Body: { shopId?: string }
  327. */
  328. router.post('/sorftime-schedule/collect-shop-products', async (req, res) => {
  329. const { shopId } = req.body || {};
  330. const { send, close } = createSseResponse(res);
  331. req.on('close', () => close());
  332. try {
  333. send('progress', { type: 'info', message: `开始采集店铺产品数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  334. const result = await sorftimeScheduler.collectShopProducts(
  335. shopId || undefined,
  336. (event) => send('progress', event)
  337. );
  338. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  339. } catch (error) {
  340. send('error', { success: false, message: error.message || '采集店铺产品时发生未知错误' });
  341. } finally {
  342. close();
  343. }
  344. });
  345. router.get('/test', (req, res)=> {
  346. res.json({
  347. message: "测试 sp-api-schedule",
  348. body: req.body
  349. });
  350. })
  351. export default router;