| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402 |
- import express from "npm:express";
- import { spApiScheduler } from './sp-api-schedule.ts';
- import { newSpApiScheduler } from './new-sp-api-schedule.ts';
- import { sorftimeScheduler } from './sorftime-api-schedule.ts';
- /**
- * 定时器清洗数据
- */
- (async () => {
- try {
- await spApiScheduler.start();
- console.log('[API Routes] Amazon SP-API定时任务已启动');
- } catch (error) {
- console.error('[API Routes] Amazon SP-API定时任务启动失败:', error);
- }
- try {
- await sorftimeScheduler.start();
- console.log('[API Routes] Sorftime定时任务已启动');
- } catch (error) {
- console.error('[API Routes] Sorftime定时任务启动失败:', error);
- }
- })();
- const router = express.Router();
- console.log('加载schedule路由 /api/schedule/');
- router.use(express.json({
- charset: 'utf-8',
- // 额外配置:确保解析URL编码的请求体也支持中文
- type: 'application/json'
- }));
- // 定时任务管理路由
- router.get('/sp-api-schedule/status', (req, res) => {
- res.json({
- success: true,
- data: {
- spApi: spApiScheduler.getStatus()
- }
- });
- });
- router.get('/sp-api-schedule/start', async (req, res) => {
- console.log('[API Routes] 启动定时任务');
- try {
- const result = await spApiScheduler.start()
- res.json({
- success: true,
- message: '定时器执行陈功',
- data: result
- });
- } catch (error) {
- res.status(500).json({
- success: false,
- message: error.message
- });
- }
-
-
- });
- router.post('/sp-api-schedule/trigger', async (req, res) => {
- try {
- const shopId = req.body.shopId;
- const result = await spApiScheduler.triggerManually(shopId);
- res.json({
- success: result.success,
- message: result.message,
- data: result
- });
- } catch (error) {
- res.status(500).json({
- success: false,
- message: error.message
- });
- }
- });
- // ===== 以下为新版独立采集接口(SSE 实时进度推送)=====
- // 前端通过 EventSource 或 fetch + ReadableStream 消费 SSE 流
- // 事件格式:
- // event: progress data: { type, message, processed? }
- // event: done data: { success, message, processed, errorCount, errors }
- // event: error data: { success: false, message }
- // 心跳注释行(': heartbeat')每 15 秒发送一次,防止连接超时
- /**
- * 创建 SSE 响应并返回辅助函数
- */
- function createSseResponse(res) {
- res.setHeader('Content-Type', 'text/event-stream; charset=utf-8');
- res.setHeader('Cache-Control', 'no-cache');
- res.setHeader('Connection', 'keep-alive');
- res.setHeader('X-Accel-Buffering', 'no');
- res.setHeader('Access-Control-Allow-Origin', '*');
- res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
- res.setHeader('Access-Control-Allow-Headers', 'Content-Type');
- res.flushHeaders();
- const send = (event, data) => {
- if (!res.writableEnded) {
- res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
- }
- };
- const heartbeatTimer = setInterval(() => {
- if (!res.writableEnded) res.write(': heartbeat\n\n');
- }, 15000);
- const close = () => {
- clearInterval(heartbeatTimer);
- if (!res.writableEnded) res.end();
- };
- return { send, close };
- }
- /**
- * POST /api/schedule/sp-api-schedule/collect-listing
- * 采集 Listing 数据,通过 SSE 流式返回进度
- * Body: { shopId: string }
- */
- router.post('/sp-api-schedule/collect-listing', async (req, res) => {
- const { shopId } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始采集 Listing 数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
- const result = await newSpApiScheduler.collectListings(
- shopId || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集 Listing 数据时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sp-api-schedule/collect-order
- * 采集订单数据,通过 SSE 流式返回进度
- * Body: { shopId: string }
- */
- router.post('/sp-api-schedule/collect-order', async (req, res) => {
- const { shopId } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始采集订单数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
- const result = await newSpApiScheduler.collectOrders(
- shopId || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集订单数据时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sp-api-schedule/collect-report
- * 采集报表数据(创建报表任务),通过 SSE 流式返回进度
- * Body: { shopId: string }
- */
- router.post('/sp-api-schedule/collect-report', async (req, res) => {
- const { shopId } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始采集报表数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
- const result = await newSpApiScheduler.collectReports(
- shopId || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集报表数据时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sp-api-schedule/parse-report
- * 检查未完成报表状态并解析已完成报表,通过 SSE 流式返回进度
- * Body: { shopId: string }
- */
- router.post('/sp-api-schedule/parse-report', async (req, res) => {
- const { shopId } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始解析报表${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
- const result = await newSpApiScheduler.parseReports(
- shopId || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '解析报表时发生未知错误' });
- } finally {
- close();
- }
- });
- // Sorftime 定时任务管理路由
- router.get('/sorftime-schedule/status', (req, res) => {
- res.json({
- success: true,
- data: {
- sorftime: sorftimeScheduler.getStatus()
- }
- });
- });
- /**
- * POST /api/schedule/sorftime-schedule/trigger-products
- * 手动触发全量类目热销产品采集(每类目分4页请求),SSE 实时返回进度
- * 对应定时任务:每月1日 00:00(Asia/Shanghai)
- * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
- * SSE done: { success, message, successCount, errorCount, errors }
- */
- router.post('/sorftime-schedule/trigger-products', async (req, res) => {
- const { shopId, nodeIds } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: '开始采集类目热销产品...' });
- const result = await sorftimeScheduler.triggerCategoryProducts(
- shopId || undefined,
- nodeIds || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集类目热销产品时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sorftime-schedule/trigger-keywords
- * 手动触发全量类目关键词采集,SSE 实时返回进度
- * 对应定时任务:每周一 01:00(Asia/Shanghai)
- * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
- * SSE done: { success, message, successCount, errorCount, errors }
- */
- router.post('/sorftime-schedule/trigger-keywords', async (req, res) => {
- const { shopId, nodeIds } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: '开始采集类目关键词...' });
- const result = await sorftimeScheduler.triggerCategoryKeywords(
- shopId || undefined,
- nodeIds || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集类目关键词时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sorftime-schedule/trigger-category-trends
- * 手动触发全量类目市场趋势采集(每类目6种 TrendIndex),SSE 实时返回进度
- * 对应定时任务:每月1日 01:00(Asia/Shanghai)
- * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
- * SSE done: { success, message, successCount, errorCount, errors }
- */
- router.post('/sorftime-schedule/trigger-category-trends', async (req, res) => {
- const { shopId, nodeIds } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: '开始采集类目市场趋势...' });
- const result = await sorftimeScheduler.triggerCategoryTrends(
- shopId || undefined,
- nodeIds || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集类目市场趋势时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sorftime-schedule/trigger-single-category
- * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势),SSE 实时返回进度
- * Body: { nodeId: string, domain?: number } domain 默认 1(美国站)
- * SSE done: { success, message, tasks: { products, keywords, trends }, duration, errors }
- */
- router.post('/sorftime-schedule/trigger-single-category', async (req, res) => {
- const { nodeId, domain } = req.body || {};
- if (!nodeId) {
- return res.status(400).json({ success: false, message: '缺少 nodeId 参数' });
- }
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始采集类目 ${nodeId} 全量数据(热销产品 + 关键词 + 市场趋势)` });
- const result = await sorftimeScheduler.triggerSingleCategory(
- nodeId,
- domain || 1,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, tasks: result.tasks, duration: result.duration, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集单类目数据时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sorftime-schedule/collect-reviews
- * 手动触发产品评论采集(遍历 Product 表 ASIN 请求 ProductReviewsQuery),SSE 实时返回进度
- * Body: 无
- * SSE done: { success, message, successCount, errorCount, errors }
- */
- router.post('/sorftime-schedule/collect-reviews', async (req, res) => {
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: '开始采集产品评论数据...' });
- const result = await sorftimeScheduler.triggerProductReviews(
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集产品评论时发生未知错误' });
- } finally {
- close();
- }
- });
- /**
- * POST /api/schedule/sorftime-schedule/collect-shop-products
- * 采集店铺竞品数据(通过 SellerId QueryType=5 分页请求),SSE 实时返回进度
- * Body: { shopId?: string }
- */
- router.post('/sorftime-schedule/collect-shop-products', async (req, res) => {
- const { shopId } = req.body || {};
- const { send, close } = createSseResponse(res);
- req.on('close', () => close());
- try {
- send('progress', { type: 'info', message: `开始采集店铺产品数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
- const result = await sorftimeScheduler.collectShopProducts(
- shopId || undefined,
- (event) => send('progress', event)
- );
- send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
- } catch (error) {
- send('error', { success: false, message: error.message || '采集店铺产品时发生未知错误' });
- } finally {
- close();
- }
- });
- router.get('/test', (req, res)=> {
- res.json({
- message: "测试 sp-api-schedule",
- body: req.body
- });
- })
- export default router;
|