feat(可回溯): 可回溯

增加文件上传
This commit is contained in:
BrandWang
2023-03-29 16:28:38 +08:00
parent b328605924
commit bf028d5034
@@ -2,36 +2,41 @@ package com.wabestway.recall.web;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONArray;
import com.wabestway.commons.http.ResObj;
import com.wabestway.engine.api.dfs.UploadFileDTO;
import com.wabestway.engine.api.dfs.UploadRespVO;
import com.wabestway.engine.api.feign.DfsStorageFeignClient;
import com.wabestway.recall.model.Order;
import com.wabestway.recall.model.OrderList;
import com.wabestway.recall.model.Record;
import com.wabestway.recall.model.ResObj;
import com.wabestway.recall.repository.OrderRepository;
import com.wabestway.recall.repository.RecordRepository;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.domain.*;
import org.springframework.web.bind.annotation.*;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.io.*;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
@Slf4j
@RestController
@RequestMapping("/order")
public class OrderController {
private final RecordRepository recordRepository;
private final OrderRepository orderRepository;
public OrderController(RecordRepository recordRepository, OrderRepository orderRepository) {
private final DfsStorageFeignClient storageFeignClient;
public OrderController(RecordRepository recordRepository, OrderRepository orderRepository, DfsStorageFeignClient storageFeignClient) {
this.recordRepository = recordRepository;
this.orderRepository = orderRepository;
this.storageFeignClient = storageFeignClient;
}
@GetMapping
@@ -85,14 +90,14 @@ public class OrderController {
Example<Order> ep = Example.of(order, matcher);
Sort s = Sort.by("createAt");
Flux<Order> flux = orderRepository.findAll(ep, s);
flux.subscribe(fo -> {
flux.filter(fo -> fo.getFileId() == null).subscribe(fo -> {
Flux<Record> rs = recordRepository.findAllByTraceIdOrderByCreateAtAsc(fo.getTraceId());
Mono<List<List<String>>> es = rs.map(rr -> rr.getEvents()).collectList();
Mono<List<String>> events = es.flatMapIterable(lists -> lists).flatMapIterable(list -> list).collectList();
events.subscribe(e -> {
JSONArray array = JSONArray.parseArray(JSON.toJSONString(e));
try {
createVideo(fo.getOrderId(), array.toJSONString());
createVideo(fo, array.toJSONString());
} catch (IOException ex) {
ex.printStackTrace();
}
@@ -101,7 +106,7 @@ public class OrderController {
return Mono.just(ResObj.ok());
}
void createVideo(String orderId, String jsonData) throws IOException {
void createVideo(Order fo, String jsonData) throws IOException {
// 创建临时目录
Path tempDir = null;
@@ -135,7 +140,23 @@ public class OrderController {
// 等待命令执行完成
int exitCode = process.waitFor();
System.out.println("Command exited with code " + exitCode);
File file = new File(tempDir.toFile(), "video.mp4");
if (file.exists()) {
byte[] fileBytes = Files.readAllBytes(file.toPath());
String titles = fo.getOrderId() + ".mp4";
UploadFileDTO uploadFileDTO = new UploadFileDTO(fileBytes, "contractExcel", titles);
ResObj<UploadRespVO> uploadResObj = storageFeignClient.uploadBytes(uploadFileDTO);
log.info(JSON.toJSONString(uploadResObj));
if (uploadResObj.isOk()) {
fo.setFileId(uploadResObj.getData().getFileId());
fo.setFileUrl(uploadResObj.getData().getFileUrl());
orderRepository.save(fo).subscribe();
}
}
} catch (IOException | InterruptedException e) {
e.printStackTrace();
} finally {