guyue
13 小时以前 18b282e5902d286f8c93dbeee19f727698c59e5c
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
package com.linghu.controller;
 
import java.time.Duration;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
 
import javax.servlet.http.HttpServletRequest;
 
import com.fasterxml.jackson.core.JsonProcessingException;
import com.linghu.model.dto.*;
import com.linghu.model.entity.*;
import com.linghu.service.*;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.ParameterizedTypeReference;
import org.springframework.dao.DuplicateKeyException;
import org.springframework.http.*;
import org.springframework.web.reactive.function.client.WebClient;
 
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.core.conditions.update.LambdaUpdateWrapper;
import com.linghu.model.common.ResponseResult;
import com.linghu.utils.JwtUtils;
 
import io.jsonwebtoken.lang.Collections;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
 
import org.springframework.web.bind.annotation.* ;
import org.springframework.http.HttpStatus;
import com.linghu.model.dto.TaskResultResponse.QuestionResult;
import com.linghu.model.dto.TaskResultResponse.UserResult;
import reactor.core.scheduler.Schedulers;
 
@RestController
@RequestMapping("/collect")
@Api(value = "采集接口", tags = "采集管理")
@Slf4j
public class CollectController {
 
    @Autowired
    private ReferenceService referenceService;
 
    @Value("${linghu.url}")
    private String baseUrl;
 
    @Autowired
    private WebClient webClient;
 
    @Autowired
    private JwtUtils jwtUtils;
    @Autowired
    private KeywordService keywordService;
    @Autowired
    private QuestionService questionService;
    @Autowired
    private KeywordTaskService keywordTaskService;
    @Autowired
    private PlatformService platformService;
    @Autowired
    private TypeService typeService;
    @Autowired
    private UserService userService;
    @Autowired
    private OrderService orderService;
    // 1. 使用线程安全的队列实现
    private final BlockingQueue<SearchTaskRequest> taskQueue = new LinkedBlockingQueue<>();
    // 添加队列访问锁
    private final ReentrantLock queueLock = new ReentrantLock();
    private static boolean isProcessing = false;
 
    @PostMapping("/search")
    @ApiOperation(value = "开始采集")
    public Mono<ResponseResult<?>> createSearchTask(
            @RequestBody SearchTaskRequest searchTaskRequest,
            HttpServletRequest request) throws JsonProcessingException {
 
        // 首先检查服务器资源
        return getServerResource()
                .flatMap(resourceResponse -> {
                    double cpuUsage = parseUsage(resourceResponse.getCpu_usage_percent());
                    double memoryUsage = parseUsage(resourceResponse.getMemory_usage_percent());
 
                    if (cpuUsage >= 90.0 || memoryUsage >= 90.0) {
 
 
                            String errorMsg = String.format("服务器资源不足,请稍后再试");
 
                        log.warn(errorMsg);
                        return Mono.just(ResponseResult.error(503, errorMsg));
                    }
 
                    // 将新的任务请求加入队列
//                    taskQueue.add(searchTaskRequest);
                    queueLock.lock();
                    try {
                        taskQueue.add(searchTaskRequest);
                    } finally {
                        queueLock.unlock();
                    }
 
                    // 如果当前没有任务在处理中,则启动任务队列的处理
                    if (!isProcessing) {
                        processNextTaskInQueue();
                    }
 
                    // 返回响应,通知用户任务已开始
                    return Mono.just(ResponseResult.success("任务已加入队列,正在处理..."));
                })
                .onErrorResume(e -> {
                    log.error("检查服务器资源失败: {}", e.getMessage(), e);
                    return Mono.just(ResponseResult.error("检查服务器资源失败: " + e.getMessage()));
                });
    }
 
    private void processNextTaskInQueue() {
        // 设置为正在处理
        isProcessing = true;
 
        // 从队列中取出下一个任务
        SearchTaskRequest nextTaskRequest = taskQueue.poll();
 
        if (nextTaskRequest != null) {
            // 处理任务
            executeBatchTask(nextTaskRequest)
                    .doFinally(signal -> {
                        // 完成后,继续处理下一个任务
                        isProcessing = false;
                        if (!taskQueue.isEmpty()) {
                            processNextTaskInQueue();  // 继续处理队列中的下一个任务
                        }
                    })
                    .subscribe();
        }
    }
 
    private Mono<ResponseResult<String>> executeBatchTask(SearchTaskRequest searchTaskRequest) {
        log.info("开始处理任务:{}", searchTaskRequest);
        log.info("++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++");
        Integer keywordId = searchTaskRequest.getKeyword_id();
        int maxConcurrentUsers = searchTaskRequest.getConfig() != null ?
                searchTaskRequest.getConfig().getMax_concurrent_users() : 3;
        List<List<UserDto>> userBatches = splitUsersIntoBatches(searchTaskRequest.getUsers(), maxConcurrentUsers, keywordId);
 
        Queue<List<UserDto>> batchQueue = new LinkedList<>(userBatches); // 用队列存储批次
 
        return Mono.just(ResponseResult.success("第一个批次已开始"))
                .doOnTerminate(() -> {
                    // 启动后台任务,继续处理批次
                    executeBatchTask(batchQueue, searchTaskRequest, keywordId)
                            .subscribe();  // 使用subscribe()启动后台任务
                });
    }
 
    private Mono<ResponseResult<?>> executeBatchTask(Queue<List<UserDto>> batchQueue, SearchTaskRequest searchTaskRequest, Integer keywordId) {
        // 如果队列为空,说明所有批次已经完成
        if (batchQueue.isEmpty()) {
            return Mono.just(ResponseResult.success("所有批次已完成"));
        }
 
        List<UserDto> currentBatch = batchQueue.poll(); // 从队列中获取当前批次
        SearchTaskRequest batchRequest = new SearchTaskRequest();
        batchRequest.setUsers(currentBatch);
        batchRequest.setQuestions(searchTaskRequest.getQuestions());
        batchRequest.setConfig(searchTaskRequest.getConfig());
        batchRequest.setSave_to_database(searchTaskRequest.getSave_to_database());
        batchRequest.setWebhook_url(searchTaskRequest.getWebhook_url());
        batchRequest.setKeyword_id(keywordId);
 
        return createSingleBatchTask(batchRequest)
                .flatMap(taskResponse -> {
//                    if (taskResponse != null && taskResponse.getTask_id() != null) {
//                        // 保存任务关联到数据库
//                        return saveKeywordTasks(keywordId, taskResponse)
//                                .then(waitForTaskCompletion(taskResponse.getTask_id(), batchQueue, searchTaskRequest, keywordId));
//                    } else {
//                        return Mono.just(ResponseResult.error("创建批次任务失败"));
//                    }
                    if (taskResponse != null && taskResponse.getTask_id() != null) {
                        // 直接等待任务完成,不再保存任务关联信息
                        return waitForTaskCompletion(taskResponse.getTask_id(), batchQueue, searchTaskRequest, keywordId);
                    } else {
                        return Mono.just(ResponseResult.error("创建批次任务失败"));
                    }
                });
    }
 
//    private Mono<Void> saveKeywordTasks(Integer keywordId, SearchTaskResponse taskResponse) {
//        if (taskResponse == null || taskResponse.getTask_id() == null) {
//            return Mono.error(new RuntimeException("任务响应无效或任务ID为空"));
//        }
//
//        KeywordTask keywordTask = new KeywordTask();
//        keywordTask.setKeyword_id(keywordId);
//        keywordTask.setTask_id(taskResponse.getTask_id());
//        keywordTask.setStatus("pending");
//
//        // 将 MyBatis-Plus 的同步方法包装为 Mono<Void>
//        return Mono.fromRunnable(() -> {
//                    boolean success = keywordTaskService.saveOrUpdate(keywordTask);
//                    if (!success) {
//                        throw new RuntimeException("保存关键词任务关联失败");
//                    }
//                })
//                .doFinally(signalType -> log.info("成功保存关键词任务关联: Task ID {}", taskResponse.getTask_id()))
//                .then();
//    }
    private Mono<ResponseResult<?>> waitForTaskCompletion(String taskId, Queue<List<UserDto>> batchQueue, SearchTaskRequest searchTaskRequest, Integer keywordId) {
        // 查询任务状态
        return getTaskStatus(taskId)
                .flatMap(statusResponse -> {
                    // 如果任务状态是"submitted"或"running",继续轮询
                    if (!"completed".equalsIgnoreCase(statusResponse.getStatus()) && !"failed".equalsIgnoreCase(statusResponse.getStatus()) && !"cancelled".equalsIgnoreCase(statusResponse.getStatus()) ) {
                        return Mono.delay(Duration.ofSeconds(5))  // 延迟 5 秒后再次查询
                                .flatMap(aLong -> waitForTaskCompletion(taskId, batchQueue, searchTaskRequest, keywordId));  // 递归调用继续等待
                    } else {
                        // 如果状态为其他状态,则继续处理下一个批次
                        return executeBatchTask(batchQueue, searchTaskRequest, keywordId);
                    }
                })
                .onErrorResume(e -> {
                    // 处理查询任务状态时的错误
                    return Mono.just(ResponseResult.error("查询任务状态失败: " + e.getMessage()));
                });
    }
    @ApiOperation(value = "查询任务状态")
    @GetMapping("/status")
    public Mono<TaskStatusResponse> getTaskStatus(String taskId) {
        return webClient.get()
                .uri(baseUrl + "/api/v1/tasks/" + taskId)
                .accept(MediaType.APPLICATION_JSON)
                .retrieve()
                .onStatus(HttpStatus::isError, response -> response.bodyToMono(TaskStatusResponse.class)
                        .flatMap(errorBody -> Mono.error(new RuntimeException(errorBody.getDetail()))))
                .bodyToMono(TaskStatusResponse.class)
                .onErrorResume(e -> {
                    // 处理错误,创建一个自定义的错误响应对象
                    TaskStatusResponse errorResponse = new TaskStatusResponse();
                    errorResponse.setStatus("ERROR");
                    errorResponse.setMessage(e.getMessage());
                    errorResponse.setDetail(e.getMessage());
                    return Mono.just(errorResponse);
                });
    }
 
 
 
   // 添加一个辅助方法来安全地将字符串转换为double
   private double parseUsage(String usageStr) {
       try {
           if (usageStr != null) {
               // 移除可能存在的百分号
               usageStr = usageStr.replace("%", "").trim();
               return Double.parseDouble(usageStr);
           }
           return 0.0;
       } catch (NumberFormatException e) {
           log.error("解析资源使用率失败: {}", e.getMessage());
           return 0.0;
       }
   }
   /*@PostMapping("/search")
   @ApiOperation(value = "开始采集")
   public Mono<ResponseResult<?>> createSearchTask(
           @RequestBody SearchTaskRequest searchTaskRequest,
           HttpServletRequest request) throws JsonProcessingException {
 
       // 首先检查服务器资源
       return getServerResource()
               .flatMap(resourceResponse -> {
                   // 将字符串类型的使用率转换为double类型
                   double cpuUsage = parseUsage(resourceResponse.getCpu_usage_percent());
                   double memoryUsage = parseUsage(resourceResponse.getMemory_usage_percent());
                   // 检查CPU和内存使用率
                   if (cpuUsage >= 90.0 || memoryUsage >= 90.0) {
                       String errorMsg = String.format("服务器资源不足:CPU使用率 %.1f%%,内存使用率 %.1f%%",
                               resourceResponse.getCpu_usage_percent(), resourceResponse.getMemory_usage_percent());
                       log.warn(errorMsg);
                       return Mono.just(ResponseResult.error(503, errorMsg));
                   }
                   Integer keywordId = searchTaskRequest.getKeyword_id();
 
                   int maxConcurrentUsers = searchTaskRequest.getConfig() != null ?
                           searchTaskRequest.getConfig().getMax_concurrent_users() : 3;
                   List<List<UserDto>> userBatches = splitUsersIntoBatches(searchTaskRequest.getUsers(), maxConcurrentUsers,keywordId);
 
                   return Flux.fromIterable(userBatches)
                           .flatMap(batch -> {
                               SearchTaskRequest batchRequest = new SearchTaskRequest();
                               batchRequest.setUsers(batch);
                               batchRequest.setQuestions(searchTaskRequest.getQuestions());
                               batchRequest.setConfig(searchTaskRequest.getConfig());
                               batchRequest.setSave_to_database(searchTaskRequest.getSave_to_database());
                               batchRequest.setWebhook_url(searchTaskRequest.getWebhook_url());
                               batchRequest.setKeyword_id(keywordId);
 
                               return createSingleBatchTask(batchRequest)
                                       .delaySubscription(Duration.ofSeconds(2)); // 批次之间添加延迟
                           }, 1) // 限制并发数为1,确保顺序执行
                           .collectList() // 收集所有批次的响应
                           .flatMap(responses ->
                                   saveKeywordTasks(keywordId, responses) // 保存关联关系
                                           .thenReturn(responses) // 返回原始响应
                           )
                           .map(responses -> ResponseResult.success(responses)) // 使用ResponseResult包装结果
                           .onErrorResume(e -> {
                               log.error("创建搜索任务失败: {}", e.getMessage(), e);
                               return Mono.just(ResponseResult.error("创建搜索任务失败: " + e.getMessage()));
                           });
               })
               .onErrorResume(e -> {
                   log.error("检查服务器资源失败: {}", e.getMessage(), e);
                   return Mono.just(ResponseResult.error("检查服务器资源失败: " + e.getMessage()));
               });
   }*/
 
//    private Mono<Void> saveKeywordTasks(Integer keywordId, List<SearchTaskResponse> taskResponses) {
//        List<KeywordTask> keywordTasks = taskResponses.stream()
//                .filter(response -> response.getTask_id() != null)
//                .map(response -> {
//                    KeywordTask keywordTask = new KeywordTask();
//                    keywordTask.setKeyword_id(keywordId);
//                    keywordTask.setTask_id(response.getTask_id());
//                    keywordTask.setStatus("pending");
//                    return keywordTask;
//                })
//                .collect(Collectors.toList());
//
//        // 将 MyBatis-Plus 的同步方法包装为 Mono<Void>
//        return Mono.fromRunnable(() -> {
//                    boolean success = keywordTaskService.saveOrUpdateBatch(keywordTasks);
//                    if (!success) {
////                        throw new RuntimeException("保存关键词任务关联失败");
//                        // 添加异常处理
//                        Mono.error( new RuntimeException("保存关键词任务关联失败"));
//                    }
//                })
//                .doFinally(signalType -> log.info("成功保存 {} 个关键词任务关联", keywordTasks.size()))
//                .then();
//    }
 
    private List<List<UserDto>> splitUsersIntoBatches(List<UserDto> users, int batchSize,Integer keywordId) {
 
        Keyword keyword = keywordService.getById(keywordId);
        if (null==keyword.getNum()){
            keyword.setNum(0);
        }
        keyword.setNum(keyword.getNum()+1);
        keywordService.updateById(keyword);
 
        List<List<UserDto>> batches = new ArrayList<>();
        for (int i = 0; i < users.size(); i += batchSize) {
            batches.add(users.subList(i, Math.min(i + batchSize, users.size())));
 
        }
        for (int i = 0; i < batches.size(); i++){
            // 创建 KeywordTask 关联,task_id 设置为 null,表示任务尚未开始
            KeywordTask keywordTask = new KeywordTask();
            keywordTask.setKeyword_id(keywordId);
            keywordTask.setTask_id(null);  // 任务ID为空
 
            keywordTask.setNum(keyword.getNum());// 任务初始状态为 pending
            keywordTaskService.save(keywordTask);  // 保存 KeywordTask
        }
 
 
        return batches;
    }
 
    private Mono<SearchTaskResponse> createSingleBatchTask(SearchTaskRequest batchRequest) {
        return webClient.post()
                .uri(baseUrl + "/api/v1/search")
                .contentType(MediaType.APPLICATION_JSON)
                .bodyValue(batchRequest)
                .retrieve()
                .onStatus(HttpStatus::is4xxClientError, response -> response.bodyToMono(String.class)
                        .flatMap(errorBody -> Mono.error(new RuntimeException(errorBody))))
                .bodyToMono(new ParameterizedTypeReference<SearchTaskResponse>() {})
                .flatMap(taskResponse -> {
                    if (taskResponse != null && taskResponse.getTask_id() != null) {
 
                        // 使用 Reactor 的方式更新数据库
                        return Mono.fromRunnable(() -> {
 
                                     //更新关键词状态
                                    LambdaUpdateWrapper<Keyword> updateWrapper = new LambdaUpdateWrapper<>();
                                    updateWrapper.eq(Keyword::getKeyword_id, batchRequest.getKeyword_id());
                                    updateWrapper.set(Keyword::getStatus, "submitted");
                                    updateWrapper.set(Keyword::getTask_id, taskResponse.getTask_id());
                                    keywordService.update(updateWrapper);
                                    //设置轮数
                                    Keyword keyword = keywordService.getById(batchRequest.getKeyword_id());
                                    // 更新关键词任务与任务ID的关联
                                    // 获取与关键词相关的任务,task_id 为 null,确保只取一个任务
                                    List<KeywordTask> keywordTasks = keywordTaskService.list(new LambdaQueryWrapper<KeywordTask>()
                                            .eq(KeywordTask::getKeyword_id, keyword.getKeyword_id())
                                            .isNull(KeywordTask::getTask_id));
                                    if (keywordTasks.size() > 0) {
                                        KeywordTask keywordTask = keywordTasks.get(0);
                                        keywordTask.setTask_id(taskResponse.getTask_id());
                                        keywordTask.setStatus("pending");
                                        keywordTaskService.updateById(keywordTask);
                                    }
//                                    KeywordTask taskToUpdate = keywordTaskService.getOne(new LambdaQueryWrapper<KeywordTask>()
//                                            .eq(KeywordTask::getKeyword_id, keyword.getKeyword_id())
//                                            .eq(KeywordTask::getTask_id, null)); // 确保 task_id 为 null 的任务
 
                                    // 更新任务的 task_id 和状态
//                                    taskToUpdate.setTask_id(taskResponse.getTask_id());
//                                    taskToUpdate.setStatus("submitted");
//                                    keywordTaskService.updateById(taskToUpdate);
 
 
                                    //设置订单进入采集状态
                                    LambdaUpdateWrapper<Orders> updateOrderWrapper = new LambdaUpdateWrapper<>();
                                    updateOrderWrapper.eq(Orders::getOrder_id, keyword.getOrder_id()) // 确保字段名正确
                                            .set(Orders::getStatus, 2); // 直接设置状态值
 
                                    boolean success = orderService.update(updateOrderWrapper);
                                    log.info("订单状态更新结果: {}", success ? "成功" : "失败");
 
                                }).subscribeOn(Schedulers.boundedElastic()) // 在弹性线程池执行
                                .thenReturn(taskResponse);
                    }
                    return Mono.just(taskResponse);
                });
    }
 
    // 移除原来的waitForTaskCompletion方法,不再需要同步等待
//    @ApiOperation(value = "查询任务状态")
//    @GetMapping("/status")
//    public Mono<TaskStatusResponse> getTaskStatus(String taskId) {
//        return webClient.get()
//                .uri(baseUrl + "/api/v1/tasks/" + taskId)
//                .accept(MediaType.APPLICATION_JSON)
//                .retrieve()
//                .onStatus(HttpStatus::isError, response -> response.bodyToMono(TaskStatusResponse.class)
//                        .flatMap(errorBody -> Mono.error(new RuntimeException(errorBody.getDetail()))))
//                .bodyToMono(TaskStatusResponse.class)
//                .flatMap(result -> {
//                    TaskStatusResponse taskStatusResponse = result;
//                    if (taskStatusResponse != null && taskStatusResponse.getStatus() != null) {
//                        List<Question> updateQuestions = taskStatusResponse.getQuestions_status().stream()
//                                .map(qs -> {
//                                    Question question = new Question();
//                                    question.setQuestion_id(qs.getQuestion_id());
//                                    question.setStatus(qs.getStatus());
//                                    return question;
//                                }).collect(Collectors.toList());
//
//                        // 包装成响应式操作
//                        return Mono.fromCallable(() -> {
//                            questionService.updateBatchById(updateQuestions);
//                            return result;
//                        });
//
//                    }
//                    return Mono.just(result);
//                })
//                .onErrorResume(e -> {
//                    // 创建一个自定义的错误响应对象
//                    TaskStatusResponse errorResponse = new TaskStatusResponse();
//                    errorResponse.setStatus("ERROR");
//                    errorResponse.setMessage(e.getMessage());
//                    errorResponse.setDetail(e.getMessage());
//
//                    return Mono.just(errorResponse);
//                });
//    }
    @PostMapping("/cancel/{keywordId}")
    @ApiOperation(value = "取消任务")
    public Mono<ResponseResult<TaskCancelResponse>> cancelTask(@PathVariable Integer keywordId) {
        // 1. 查询所有与关键词相关的任务
        List<KeywordTask> tasks = keywordTaskService.list(
                new LambdaQueryWrapper<KeywordTask>().eq(KeywordTask::getKeyword_id, keywordId)
        );
 
        // 2. 从队列中移除所有相关任务(使用锁保证线程安全)
        List<SearchTaskRequest> removedQueueTasks = removeTasksFromQueueByKeywordId(keywordId);
 
        // 3. 筛选出需要远程取消的任务
        List<KeywordTask> tasksToCancelRemotely = tasks.stream()
                .filter(task -> task.getTask_id() != null && "pending".equalsIgnoreCase(task.getStatus()))
                .collect(Collectors.toList());
 
        // 4. 对筛选出的任务发送远程取消请求(并行执行)
        return Flux.fromIterable(tasksToCancelRemotely)
                .flatMap(task -> {
                    // 创建状态更新和远程取消的组合操作
                    Mono<Void> updateStatus = updateTaskStatus(task.getTask_id(), "canceled");
                    Mono<ResponseResult<?>> cancelOp = cancelRemoteTask(task.getTask_id())
                            .onErrorResume(e -> {
                                log.error("取消任务 {} 失败: {}", task.getTask_id(), e.getMessage());
                                return Mono.just(ResponseResult.error("取消任务失败: " + e.getMessage()));
                            });
 
                    // 合并操作:无论远程取消是否成功,都更新状态
                    return Mono.zip(cancelOp, updateStatus)
                            .thenReturn(true);
                }, 10) // 设置10的并发度
                .collectList()
                .thenReturn(ResponseResult.success(
                        new TaskCancelResponse(
                                String.format("关键词任务已取消,队列中移除 %d 个任务,远程取消 %d 个任务",
                                        removedQueueTasks.size(),
                                        tasksToCancelRemotely.size())
                        )
                ))
                .onErrorResume(e -> {
                    log.error("取消关键词任务失败: {}", e.getMessage());
                    return Mono.just(ResponseResult.error(500, "取消关键词任务失败: " + e.getMessage()));
                });
    }
 
//    @PostMapping("/cancel/{keywordId}")
//    @ApiOperation(value = "取消任务")
//    public Mono<ResponseResult<TaskCancelResponse>> cancelTask(@PathVariable Integer keywordId) {
//        // 1. 查询所有与关键词相关的任务
//        List<KeywordTask> tasks = keywordTaskService.list(
//                new LambdaQueryWrapper<KeywordTask>().eq(KeywordTask::getKeyword_id, keywordId)
//        );
//
//        // 2. 从队列中移除所有相关任务
//        List<SearchTaskRequest> removedQueueTasks = removeTasksFromQueueByKeywordId(keywordId);
//
//        // 3. 筛选出需要远程取消的任务(任务ID不为空且状态为pending)
//        List<KeywordTask> tasksToCancelRemotely = tasks.stream()
//                .filter(task -> task.getTask_id() != null && "pending".equalsIgnoreCase(task.getStatus()))
//                .collect(Collectors.toList());
//        // 检查是否有任务与关键词相关
////        if (tasks.isEmpty()) {
////            return Mono.just(ResponseResult.error("没有找到相关任务"));
////        }
//
//        // 4. 对筛选出的任务发送远程取消请求
//        List<Mono<ResponseResult<?>>> cancelRequests = tasksToCancelRemotely.stream()
//                .map(task -> cancelRemoteTask(task.getTask_id())
//                        .doOnSuccess(response -> {
//                            // 更新任务状态为canceled
//                            updateTaskStatus(task.getTask_id(), "canceled").subscribe();
//                        })
//                        .onErrorResume(e -> {
//                            log.error("取消任务 {} 失败: {}", task.getTask_id(), e.getMessage());
//                            // 即使取消失败,也尝试更新状态
//                            updateTaskStatus(task.getTask_id(), "canceled").subscribe();
//                            return Mono.just(ResponseResult.error("取消任务失败: " + e.getMessage()));
//                        }))
//                .collect(Collectors.toList());
//
//        // 5. 并行执行所有取消请求
//        return Flux.fromIterable(cancelRequests)
//                .concatMap(request -> request) // 顺序执行,而非并行
//                .collectList()
//                .thenReturn(ResponseResult.success(
//                        new TaskCancelResponse(
//                                String.format("关键词任务已取消,队列中移除 %d 个任务,远程取消 %d 个任务",
//                                        removedQueueTasks.size(),
//                                        tasksToCancelRemotely.size())
//                        )
//                ))
//                .onErrorResume(e -> {
//                    log.error("取消关键词任务失败: {}", e.getMessage());
//                    return Mono.just(ResponseResult.error(500, "取消关键词任务失败: " + e.getMessage()));
//                });
//    }
// 线程安全的队列移除方法
private List<SearchTaskRequest> removeTasksFromQueueByKeywordId(Integer keywordId) {
    List<SearchTaskRequest> removedTasks = new ArrayList<>();
 
    // 使用锁保证队列操作的原子性
    queueLock.lock();
    try {
        Iterator<SearchTaskRequest> iterator = taskQueue.iterator();
        while (iterator.hasNext()) {
            SearchTaskRequest task = iterator.next();
            if (task.getKeyword_id() != null && task.getKeyword_id().equals(keywordId)) {
                removedTasks.add(task);
                iterator.remove();
            }
        }
    } finally {
        queueLock.unlock();
    }
 
    log.info("从队列中移除了 {} 个与关键词ID {} 相关的任务", removedTasks.size(), keywordId);
    return removedTasks;
}
 
    // 从队列中移除所有关键词ID匹配的任务
//    private List<SearchTaskRequest> removeTasksFromQueueByKeywordId(Integer keywordId) {
//        List<SearchTaskRequest> removedTasks = new ArrayList<>();
//
//        // 使用迭代器安全地移除元素
//        Iterator<SearchTaskRequest> iterator = taskQueue.iterator();
//        while (iterator.hasNext()) {
//            SearchTaskRequest task = iterator.next();
//            if (task.getKeyword_id() != null && task.getKeyword_id().equals(keywordId)) {
//                removedTasks.add(task);
//                iterator.remove();
//            }
//        }
//
//        log.info("从队列中移除了 {} 个与关键词ID {} 相关的任务", removedTasks.size(), keywordId);
//        return removedTasks;
//    }
    // 发送远程取消请求
// 发送远程取消请求(使用Java 8兼容的Map创建方式)
    private Mono<ResponseResult<?>> cancelRemoteTask(String taskId) {
        // 使用Collections.singletonMap或手动创建Map
        Map<String, Object> requestBody = new HashMap<>();
        requestBody.put("status", "pending");
 
        return webClient.post()
                .uri(baseUrl + "/api/v1/tasks/" + taskId + "/cancel")
                .contentType(MediaType.APPLICATION_JSON)
                .bodyValue(requestBody)
                .retrieve()
                .onStatus(HttpStatus::isError, response -> response.bodyToMono(String.class)
                        .flatMap(errorBody -> Mono.error(new RuntimeException("取消失败: " + errorBody))))
                .bodyToMono(Void.class)
                .thenReturn(ResponseResult.success("任务已取消"));
    }
 
    // 更新单个任务状态
    private Mono<Void> updateTaskStatus(String taskId, String status) {
        return Mono.fromRunnable(() -> {
            LambdaUpdateWrapper<KeywordTask> updateWrapper = new LambdaUpdateWrapper<>();
            updateWrapper.eq(KeywordTask::getTask_id, taskId);
            updateWrapper.set(KeywordTask::getStatus, status);
            keywordTaskService.update(updateWrapper);
        }).subscribeOn(Schedulers.boundedElastic()).then();
    }
 /*   @PostMapping("/cancel/{taskId}")
    @ApiOperation(value = "取消任务")
    public Mono<ResponseResult<TaskCancelResponse>> cancelTask(@PathVariable String taskId) {
        return webClient.post()
                .uri(baseUrl + "/api/v1/tasks/" + taskId + "/cancel")
                .contentType(MediaType.APPLICATION_JSON)
                .bodyValue(Collections.emptyMap()) // 添加空请求体
                .retrieve()
                .onStatus(HttpStatus::isError, response -> response.bodyToMono(TaskCancelResponse.class)
                        .flatMap(errorBody -> Mono.error(new RuntimeException(errorBody.getDetail()))))
                .bodyToMono(TaskCancelResponse.class)
                .flatMap(cancelResponse -> {
                    // 更新关键词状态
                    Mono<Void> updateKeyword = Mono.fromRunnable(() -> {
                                LambdaUpdateWrapper<Keyword> updateWrapper = new LambdaUpdateWrapper<>();
                                updateWrapper.eq(Keyword::getTask_id, taskId);
                                updateWrapper.set(Keyword::getStatus, "canceled"); // 统一使用"canceled"
                                keywordService.update(updateWrapper);
                            })
                            .subscribeOn(Schedulers.boundedElastic())
                            .then();
 
                    // 更新关键词任务状态
                    Mono<Void> updateKeywordTask = Mono.fromRunnable(() -> {
                                LambdaUpdateWrapper<KeywordTask> updateWrapper = new LambdaUpdateWrapper<>();
                                updateWrapper.eq(KeywordTask::getTask_id, taskId);
                                updateWrapper.set(KeywordTask::getStatus, "canceled"); // 统一使用"canceled"
                                keywordTaskService.update(updateWrapper);
                            })
                            .subscribeOn(Schedulers.boundedElastic())
                            .then();
 
                    // 并行执行两个更新操作,并在完成后返回cancelResponse
                    return Mono.when(updateKeyword, updateKeywordTask)
                            .thenReturn(cancelResponse);
                })
                .map(data -> ResponseResult.success(data))
                .onErrorResume(e -> {
                    if (e.getMessage().contains("任务不存在")) {
                        return Mono.just(ResponseResult.error(200, e.getMessage()));
                    } else if (e.getMessage().contains("无法取消")) {
                        return Mono.just(ResponseResult.error(200, e.getMessage()));
                    }
                    return Mono.just(ResponseResult.error(500,  e.getMessage()));
                });
    }*/
 
    @ApiOperation(value = "获取任务结果")
    @GetMapping("/tasks/{taskId}")
    public Mono<TaskResultResponse> getTaskResult(@PathVariable String taskId) {
        return webClient.get()
                .uri(baseUrl + "/api/v1/tasks/" + taskId + "/result")
                .accept(MediaType.APPLICATION_JSON)
                .retrieve()
                .onStatus(HttpStatus::is4xxClientError, response -> {
                    if (response.statusCode() == HttpStatus.NOT_FOUND) {
                        return response.bodyToMono(String.class)
                                .flatMap(errorBody -> Mono.error(new RuntimeException("任务不存在")));
                    } else if (response.statusCode() == HttpStatus.BAD_REQUEST) {
                        return response.bodyToMono(String.class)
                                .flatMap(errorBody -> Mono.error(new RuntimeException("任务未完成,无法获取结果")));
                    }
                    return response.createException().flatMap(Mono::error);
                })
                .bodyToMono(TaskResultResponse.class)
                .flatMap(responseResult -> {
                    TaskResultResponse result = responseResult;
 
                    if (result != null && result.getResults() != null) {
                        return updateQuestionAndReference(result)
                                .thenReturn(responseResult);
                    }
                    return Mono.just(responseResult);
                })
                .onErrorResume(e -> {
                    System.out.println("获取任务结果失败");
                        TaskResultResponse result = new TaskResultResponse();
                    result.setDetail("获取任务结果失败: " + e.getMessage());
                    return Mono.just(result);
                });
    }
 
//    private Mono<Void> updateQuestionAndReference(TaskResultResponse result) {
//        return Mono.fromRunnable(() -> {
//            // 1. 更新关键词状态
//            LambdaUpdateWrapper<Keyword> keywordUpdate = new LambdaUpdateWrapper<>();
//            keywordUpdate.eq(Keyword::getTask_id, result.getTask_id())
//                    .set(Keyword::getStatus, "completed");
//            keywordService.update(keywordUpdate);
//
//            // 查询关键词ID
//            LambdaQueryWrapper<Keyword> keywordQuery = new LambdaQueryWrapper<>();
//            keywordQuery.eq(Keyword::getTask_id, result.getTask_id());
//            Keyword keyword = keywordService.getOne(keywordQuery);
//
//            if (keyword == null) {
//                System.out.println("未找到关联的关键词,task_id: " + result.getTask_id());
//                return;
//            }
//
//            // 2. 处理每个用户的问题结果
//            for (UserResult userResult : result.getResults()) {
//                for (QuestionResult questionResult : userResult.getQuestions_results()) {
//                    // 2.1 查询问题ID
//                    LambdaQueryWrapper<Question> queryWrapper = new LambdaQueryWrapper<>();
//                    queryWrapper.eq(Question::getQuestion, questionResult.getQuestion())
//                            .eq(Question::getKeyword_id, keyword.getKeyword_id());
//                    Question question = questionService.getOne(queryWrapper);
//
//                    if (question != null) {
//                        // 更新问题状态
//                        LambdaUpdateWrapper<Question> updateWrapper = new LambdaUpdateWrapper<>();
//                        updateWrapper.eq(Question::getQuestion_id, question.getQuestion_id())
//                                .set(Question::getStatus, questionResult.getStatus())
//                                .set(Question::getResponse, questionResult.getResponse())
//                                .set(Question::getExtracted_count, questionResult.getExtracted_count())
//                                .set(Question::getError, questionResult.getError())
//                                .set(Question::getTimestamp, LocalDateTime.parse(
//                                        questionResult.getTimestamp(),
//                                        DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSSSS")
//                                ));
//                        questionService.update(updateWrapper);
//
//                        // 2.2 保存引用数据
//                        List<Reference> references = questionResult.getReferences().stream()
//                                .map(ref -> {
//                                    Reference reference = new Reference();
//                                    reference.setQuestion_id(question.getQuestion_id());
//                                    reference.setTitle(ref.getTitle());
//                                    reference.setUrl(ref.getUrl());
//                                    reference.setDomain(ref.getDomain());
//                                    reference.setCreate_time(LocalDateTime.now());
//                                    return reference;
//                                })
//                                .collect(Collectors.toList());
//
//                        if (!references.isEmpty()) {
//                            referenceService.saveBatch(references);
//                        }
//                    } else {
//                        System.out.println("未找到匹配的问题,question " + question.getQuestion());
//
//                    }
//                }
//            }
//        });
//    }
    /**
     * 获取或创建平台(确保同一domain只创建一次)
     * @param domain 平台域名
     * @return 已存在或新创建的Platform
     */
    private Platform getOrCreatePlatform(String domain) {
        // 1. 先尝试查询已存在的平台
        Platform platform = platformService.getPlatformByDomain(domain);
        if (platform != null) {
            return platform;
        }
 
        // 2. 若不存在,尝试创建(处理并发场景)
        try {
            // 2.1 获取或创建“默认”类型(Type也需避免重复,建议Type表的type_name也加唯一约束)
            Type defaultType = typeService.getOne(new LambdaQueryWrapper<Type>()
                    .eq(Type::getType_name, "默认"));
            if (defaultType == null) {
                defaultType = new Type();
                defaultType.setType_name("默认");
                typeService.save(defaultType); // 若Type可能重复,此处也需处理DuplicateKeyException
            }
 
            // 2.2 构建新平台对象
            Platform newPlatform = new Platform();
            newPlatform.setDomain(domain);
            newPlatform.setPlatform_name(domain); // 平台名称默认使用域名,可根据实际需求调整
            newPlatform.setType_id(defaultType.getType_id());
            newPlatform.setCreate_time(LocalDateTime.now()); // 补充创建时间
 
            // 2.3 尝试保存,若因唯一约束冲突失败,则捕获异常
            platformService.save(newPlatform);
            return newPlatform; // 保存成功,返回新创建的平台
 
        } catch (DuplicateKeyException e) {
            // 3. 若捕获到重复键异常,说明并发创建了,重新查询即可(此时数据库中已存在该平台)
            log.warn("平台domain={}已存在,无需重复创建", domain, e);
            return platformService.getPlatformByDomain(domain); // 重新查询,一定能获取到
        } catch (Exception e) {
            // 处理其他异常(如数据库连接失败等)
            log.error("创建平台失败,domain={}", domain, e);
            throw new RuntimeException("创建平台失败", e);
        }
    }
    private Mono<Void> updateQuestionAndReference(TaskResultResponse result) {
        return Mono.fromRunnable(() -> {
            try {
                //查看每个账号信息的status是否正常
 
                // 1. 根据KeywordTask更新关键词状态
                // 查询关键词ID
                LambdaQueryWrapper<KeywordTask> keywordTaskWrapper = new LambdaQueryWrapper<>();
                keywordTaskWrapper.eq(KeywordTask::getTask_id, result.getTask_id());
                KeywordTask keywordTask = keywordTaskService.getOne(keywordTaskWrapper);
//                keywordTask.setStatus("completed");
//                keywordTaskService.updateById(keywordTask);
                Keyword keyword = keywordService.getById(keywordTask.getKeyword_id());
 
                if (keyword == null) {
                    System.out.println("未找到关联的关键词,task_id: " + result.getTask_id());
                    //报错
                    throw new Exception("未找到关联的关键词,task_id: " + result.getTask_id());
//                    return;
                }
                LambdaQueryWrapper<KeywordTask> keywordTaskWrapper2 = new LambdaQueryWrapper<>();
                keywordTaskWrapper2.eq(KeywordTask::getKeyword_id, keyword.getKeyword_id());
                List<KeywordTask> keywordTasks = keywordTaskService.list(keywordTaskWrapper2);
 
                //如果全部为completed  关键词也为completed  ,如果关联关系没有任务id,或者状态为running  ,关键词为submitted,
                if (keywordTasks.stream().allMatch(task -> "completed".equals(task.getStatus())) ) {
                    keyword.setStatus("completed");
                    keywordService.updateById(keyword);
 
                }
                //如果有一个task为failed设置关键词为false
                else if (keywordTasks.stream().anyMatch(task -> "failed".equals(task.getStatus()))) {
                    keyword.setStatus("false");
                    keywordService.updateById(keyword);
                }
 
 
 
//
//// 定义状态优先级:canceled > false > completed
//                String finalStatus = "completed"; // 默认状态为 completed
//
//                for (KeywordTask task : keywordTasks) {
//                    String status = task.getStatus();
////                    if ("canceled".equals(status)) {
////                        finalStatus = "canceled";
////                        break; // 遇到 canceled 直接跳出循环,因为优先级最高
////                    } else
//                        if ("false".equals(status)) {
//                        finalStatus = "false";
//                        // 不跳出循环,继续检查是否存在 canceled
//                    }
//                }
// 更新关键词状态
 
                    String orderId = keyword.getOrder_id();
                    if (orderId == null || orderId.isEmpty()) {
                        System.out.println("关键词[" + keyword.getKeyword_id() + "]未关联订单,跳过订单状态更新");
                        return;
                    }
 
                    // 2. 查询该订单下的所有关键词
                    LambdaQueryWrapper<Keyword> orderKeywordsWrapper = new LambdaQueryWrapper<>();
                    orderKeywordsWrapper.eq(Keyword::getOrder_id, orderId);
                    List<Keyword> orderKeywords = keywordService.list(orderKeywordsWrapper);
 
                    if (orderKeywords.isEmpty()) {
                        System.out.println("订单[" + orderId + "]下无关键词,跳过状态更新");
                        return;
                    }
                    // 3. 检查所有关键词的状态是否均为 completed 或 false
                    boolean allValid = orderKeywords.stream()
                            .allMatch(k -> "completed".equals(k.getStatus()) || "false".equals(k.getStatus()) || "canceled".equals(k.getStatus()));
 
                    // 4. 若所有关键词状态均有效,更新订单状态为3
                    if (allValid) {
                        Orders orders = orderService.getById(orderId);
                        if (orders != null) {
                            orders.setStatus(3); // 假设Orders有Integer类型的status字段
                            orderService.updateById(orders);
                            System.out.println("订单[" + orderId + "]所有关键词状态符合条件,已更新状态为3");
                        } else {
                            System.out.println("未找到订单[" + orderId + "],无法更新状态");
                        }
                    }
 
                Orders orders = orderService.getById(keyword.getOrder_id());
 
                // 2. 批量查询所有问题
                LambdaQueryWrapper<Question> queryWrapper = new LambdaQueryWrapper<>();
                queryWrapper.eq(Question::getKeyword_id, keyword.getKeyword_id());
                List<Question> questions = questionService.list(queryWrapper);
 
                // 构建问题映射表,用于快速查找
                Map<String, Question> questionMap = questions.stream()
                        .collect(Collectors.toMap(Question::getQuestion, q -> q));
 
                // 3. 收集所有需要更新的问题和引用
                List<Question> questionsToUpdate = new ArrayList<>();
                List<Reference> allReferences = new ArrayList<>();
                List<Reference> resultList = new ArrayList<>();
                // 遍历结果
                for (UserResult userResult : result.getResults()) {
                    for (QuestionResult questionResult : userResult.getQuestions_results()) {
                        try {
                            Question question = questionMap.get(questionResult.getQuestion());
                            if (question != null) {
                                // 更新问题对象
                                question.setStatus(questionResult.getStatus());
                                question.setResponse(questionResult.getResponse());
                                question.setExtracted_count(questionResult.getExtracted_count());
                                question.setError(questionResult.getError());
                                question.setKeyword_id(keyword.getKeyword_id());
 
                                // 解析时间戳
                                if (questionResult.getTimestamp() != null) {
                                    DateTimeFormatter formatter = DateTimeFormatter
                                            .ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSSSS");
                                    question.setTimestamp(
                                            LocalDateTime.parse(questionResult.getTimestamp(), formatter));
                                }
                                //更新
//                                questionService.updateById(question);
 
                                questionsToUpdate.add(question);
 
                               /* List<Reference> references =
                                        Optional.ofNullable(questionResult.getReferences())
                                                .orElse(Collections.emptyList())
                                                .stream()
                                                .map(ref -> {
                                                    Reference reference = new Reference();
                                                    reference.setQuestion_id(question.getQuestion_id());
                                                    reference.setTitle(ref.getTitle());
                                                    reference.setUrl(ref.getUrl());
                                                    reference.setDomain(ref.getDomain());
                                                    reference.setNum(keyword.getNum());
                                                    reference.setTask_id(result.getTask_id());
                                                    reference.setKeyword_id(keyword.getKeyword_id());
                                                    //域名和平台id映射
                                                    reference.setCreate_time(LocalDateTime.now());
                                                    Platform platform = platformService.getPlatformByDomain(reference.getDomain());
                                                    if (platform == null) {
                                                        //平台为空 创建平台 类型为“默认”
                                                        Type type = typeService.getOne(new LambdaQueryWrapper<Type>().eq(Type::getType_name,"默认"));
                                                        if (type == null) {
                                                            Type newType = new Type();
                                                            newType.setType_name("默认");
                                                            typeService.save(newType);
                                                            type = newType;
                                                        }
                                                        Platform platform1 = new Platform();
                                                        platform1.setDomain(reference.getDomain());
                                                        platform1.setPlatform_name(reference.getDomain());
                                                        platform1.setType_id(type.getType_id());
                                                        platformService.save(platform1);
 
                                                        reference.setType_id(type.getType_id());
                                                        reference.setPlatform_id(platform1.getPlatform_id());
 
                                                    }
                                                    else {
                                                        reference.setPlatform_id(platform.getPlatform_id());
                                                        Type type = typeService.getById(platform.getType_id());
                                                        if (type != null){
                                                            reference.setType_id(type.getType_id());
                                                        }
                                                    }
                                                    return reference;
                                                })
                                                .collect(Collectors.toList());*/
                                // 初始化引用列表(避免null)
                                List<Reference> references = new ArrayList<>();
                                List<TaskResultResponse.Reference> originalReferences = questionResult.getReferences();
                                if (originalReferences == null) {
                                    originalReferences = Collections.emptyList();
                                }
 
                                // 遍历原始引用列表,转换为Reference对象
                                for (TaskResultResponse.Reference ref : originalReferences) { // 注意:需将“原引用类型”替换为实际类型(如QuestionResult中的引用类型)
                                    Reference reference = new Reference();
                                    // 设置基本字段
                                    reference.setQuestion_id(question.getQuestion_id());
                                    reference.setTitle(ref.getTitle());
                                    reference.setUrl(ref.getUrl());
                                    reference.setDomain(ref.getDomain());
                                    reference.setNum(keyword.getNum());
                                    reference.setTask_id(result.getTask_id());
                                    reference.setKeyword_id(keyword.getKeyword_id());
                                    reference.setCreate_time(LocalDateTime.now());
 
                                   /* // 处理平台和类型关联
                                    Platform platform = platformService.getPlatformByDomain(reference.getDomain());
                                    if (platform == null) {
                                        // 平台不存在,创建新平台(类型默认为“默认”)
                                        Type type = typeService.getOne(new LambdaQueryWrapper<Type>().eq(Type::getType_name, "默认"));
                                        if (type == null) {
                                            Type newType = new Type();
                                            newType.setType_name("默认");
                                            typeService.save(newType);
                                            type = newType;
                                        }
                                        Platform platform1 = new Platform();
                                        platform1.setDomain(reference.getDomain());
                                        platform1.setPlatform_name(reference.getDomain());
                                        platform1.setType_id(type.getType_id());
                                        platform1.setCreate_time(LocalDateTime.now());
                                        platformService.save(platform1);
 
                                        // 关联新平台和类型
                                        reference.setType_id(type.getType_id());
                                        reference.setPlatform_id(platform1.getPlatform_id());
                                    } else {
                                        // 平台已存在,直接关联
                                        reference.setPlatform_id(platform.getPlatform_id());
                                        Type type = typeService.getById(platform.getType_id());
                                        if (type != null) {
                                            reference.setType_id(type.getType_id());
                                        }
                                    }
*/
                                    // 关键:使用优化后的方法获取平台,避免重复创建
                                    Platform platform = getOrCreatePlatform(ref.getDomain());
                                    reference.setPlatform_id(platform.getPlatform_id());
                                    reference.setType_id(platform.getType_id()); // 直接从平台获取类型ID,更可靠
                                    // 添加到结果列表
                                    references.add(reference);
                                }
                                // 添加到总引用列表
                                if (!references.isEmpty()) {
                                    allReferences.addAll(references);
                                }
 
                                //取数据库中当前关键词的当前轮次的当前问题id结果拿出来
                                List<Reference> dbList = referenceService.list(new LambdaQueryWrapper<Reference>().eq(Reference::getKeyword_id, keyword.getKeyword_id())
                                        .eq(Reference::getNum, keyword.getNum())
                                        .eq(Reference::getQuestion_id, question.getQuestion_id())
                                );
 
                                // 1. 合并两个列表
                                List<Reference> combinedList = new ArrayList<>();
                                combinedList.addAll(allReferences);
                                combinedList.addAll(dbList);
 
                                // 2. 创建复合键的Map,用于统计完全匹配的记录
                                Map<String, List<Reference>> compositeKeyMap = combinedList.stream()
                                        .collect(Collectors.groupingBy(
                                                ref -> ref.getTitle() + "|" + ref.getUrl() + "|" + ref.getDomain()
                                        ));
 
                                // 3. 处理每组重复记录
 
                                compositeKeyMap.forEach((key, refGroup) -> {
                                    // 3.1 找出组内有ID的记录(优先从dbList中获取)
                                    Optional<Reference> existingRecord = refGroup.stream()
                                            .filter(ref -> ref.getReference_id() != null)
                                            .findFirst();
 
                                    // 3.2 统计该组的重复次数(总数-1)
                                    int repetitionCount = refGroup.size() - 1;
 
                                    // 3.3 决定最终保留的记录
                                    Reference recordToSave = new Reference();
                                    if (existingRecord.isPresent()) {
                                        // 使用已有ID的记录并更新重复次数
                                        recordToSave = existingRecord.get();
                                        recordToSave.setRepetition_num(
                                                (recordToSave.getRepetition_num() == null ? 1 : recordToSave.getRepetition_num())
                                                        + repetitionCount
                                        );
                                    } else {
                                        // 没有ID记录则取第一条并设置重复次数
                                        recordToSave = refGroup.get(0);
                                        recordToSave.setRepetition_num(1+repetitionCount);
                                    }
 
                                    resultList.add(recordToSave);
                                });
                                referenceService.saveOrUpdateBatch(resultList);
                            }
                        } catch (Exception e) {
                            log.error(e.getMessage(), e);
                            System.out.println("处理问题结果失败: " + e.getMessage());
                        }
                    }
                    //更新账号状态
                    if ( "failed".equals(userResult.getStatus())){
                        if (userResult.getError().contains("登录失败")){
                            LambdaUpdateWrapper<User> userWrapper =  new LambdaUpdateWrapper<>();
                            userWrapper.eq(User::getUser_email, userResult.getUser_email());
                            userWrapper.set(User::getStatus, "无法登录");
                            userService.update(userWrapper);
 
                        }else if (userResult.getError().contains("信息错误")){
                            LambdaUpdateWrapper<User> userWrapper =  new LambdaUpdateWrapper<>();
                            userWrapper.eq(User::getUser_email, userResult.getUser_email());
                            userWrapper.set(User::getStatus, "信息错误");
                            userService.update(userWrapper);
                        }
                    }
 
                }
 
                // 4. 批量更新问题
                System.out.println(questionsToUpdate);
                if (!questionsToUpdate.isEmpty()) {
                    questionService.updateBatchById(questionsToUpdate);
                    System.out.println("成功批量更新 " + questionsToUpdate.size() + " 个问题");
                }
 
                // 5. 批量插入引用,使用流式分批处理
//                if (!allReferences.isEmpty()) {
//                    int batchSize = 1000;
//                    IntStream.iterate(0, i -> i + batchSize)
//                            .limit((allReferences.size() + batchSize - 1) / batchSize)
//                            .forEach(i -> {
//                                List<Reference> batch = allReferences.subList(
//                                        i, Math.min(i + batchSize, allReferences.size()));
//                                referenceService.saveBatch(batch);
//                            });
//                    System.out.println("成功批量插入 " + allReferences.size() + " 条引用数据");
//                }
 
 
            } catch (Exception e) {
                log.error("更新问题和引用数据失败: " ,e.getMessage(), e);
//                System.out.println("更新问题和引用数据失败: " + e.getMessage());
                throw new RuntimeException("更新问题和引用数据失败", e);
            }
        });
    }
 
    @GetMapping("/tasks/all")
    @ApiOperation(value = "获取所有任务列表")
    public Mono<TaskListResponse> getAllTasks() {
        return webClient.get()
                .uri(baseUrl + "/api/v1/tasks")
                .accept(MediaType.APPLICATION_JSON)
                .retrieve()
                .bodyToMono(new ParameterizedTypeReference<TaskListResponse>() {
                })
                .onErrorResume(e -> {
                    TaskListResponse response = new TaskListResponse();
                    response.setDetail("获取任务列表失败: " + e.getMessage());
                    return Mono.just(response);
 
                    // return Mono.just(ResponseResult.error("获取任务列表失败: " + e.getMessage()));
                });
    }
 
    @GetMapping("/health")
    @ApiOperation("健康检查")
    public Mono<HealthResponse> checkThirdPartyHealth() {
        return webClient.get()
                .uri(baseUrl + "/health") // 假设第三方健康检查接口路径为/health
                .retrieve()
                .bodyToMono(HealthResponse.class)
                .onErrorResume(e -> Mono.just(
                        new HealthResponse("unhealthy", null, "", e.getMessage())));
    }
 
    /**
     * 查询服务器资源
     */
    @GetMapping("/server/resource")
    @ApiOperation(value = "查询服务器资源")
    public Mono<ServerResourceResponse> getServerResource() {
        return webClient.get()
                .uri(baseUrl + "/api/v1/system/resources")
                .retrieve()
                .bodyToMono(ServerResourceResponse.class)
                .onErrorResume(e -> Mono.just(
                        new ServerResourceResponse( e.getMessage())));
    }
    /**
     * 传入orderid查所有关键词id以及关键词下面的所有任务id,轮询所有任务状态,如果状态为completed,则循环调用获取结果接口,处理结果
     */
 
 
}