feat(可回溯): 重构
重构为mysql
This commit is contained in:
@@ -0,0 +1,196 @@
|
||||
package com.wabestway.recall.web;
|
||||
|
||||
import cn.hutool.core.io.FileUtil;
|
||||
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.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.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
@Slf4j
|
||||
@RestController
|
||||
@RequestMapping("/order")
|
||||
public class OrderController {
|
||||
private final RecordRepository recordRepository;
|
||||
private final OrderRepository orderRepository;
|
||||
|
||||
private final DfsStorageFeignClient storageFeignClient;
|
||||
|
||||
public OrderController(RecordRepository recordRepository, OrderRepository orderRepository, DfsStorageFeignClient storageFeignClient) {
|
||||
this.recordRepository = recordRepository;
|
||||
this.orderRepository = orderRepository;
|
||||
this.storageFeignClient = storageFeignClient;
|
||||
}
|
||||
|
||||
@GetMapping
|
||||
public Flux<Order> getAllOrders() {
|
||||
return orderRepository.findAll();
|
||||
}
|
||||
|
||||
@PostMapping("/list")
|
||||
public Mono<ResObj> list(@RequestBody Order order) {
|
||||
OrderList resList = new OrderList();
|
||||
List<Order> list = new ArrayList<>();
|
||||
Sort s = Sort.by("createAt").descending();
|
||||
PageRequest pageRequest = PageRequest.of(order.getPage() - 1, order.getPageSize(), s);
|
||||
ExampleMatcher matcher = ExampleMatcher.matching()
|
||||
.withIgnoreNullValues().withIgnorePaths("createAt", "last", "page", "pageSize");
|
||||
Example<Order> ep = Example.of(order, matcher);
|
||||
Flux<Order> flux = orderRepository.findAll(ep, s);
|
||||
flux.skip(pageRequest.getOffset()).limitRequest(pageRequest.getPageSize()).subscribe(list::add);
|
||||
resList.setList(list);
|
||||
return flux.count().flatMap(count -> {
|
||||
resList.setTotal(count);
|
||||
return Mono.just(ResObj.ok(resList));
|
||||
});
|
||||
}
|
||||
|
||||
@PostMapping("/recallUp")
|
||||
public Mono<ResObj> recallUp(@RequestBody Order order) {
|
||||
Mono<Order> queryOrder = orderRepository.findByOrderId(order.getOrderId());
|
||||
return queryOrder.flatMap(ss -> {
|
||||
ss.setStartDate(order.getStartDate());
|
||||
ss.setEndDate(order.getEndDate());
|
||||
ss.setHolderName(order.getHolderName());
|
||||
ss.setHolderPhone(order.getHolderPhone());
|
||||
ss.setSupplierName(order.getSupplierName());
|
||||
orderRepository.save(ss).subscribe();
|
||||
return Mono.just(ResObj.ok());
|
||||
}).defaultIfEmpty(ResObj.fail("订单不存在"));
|
||||
}
|
||||
|
||||
@PostMapping("/info/{orderId}")
|
||||
public Mono<ResObj> findByOrderId(@PathVariable String orderId) {
|
||||
Mono<Order> forder = orderRepository.findByOrderId(orderId);
|
||||
return forder.flatMap(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();
|
||||
return events.flatMap(e -> {
|
||||
fo.setEvents(e);
|
||||
return Mono.just(ResObj.ok(fo));
|
||||
});
|
||||
}).defaultIfEmpty(ResObj.fail("订单不存在"));
|
||||
}
|
||||
|
||||
@GetMapping("/ffmpeg")
|
||||
public Mono<ResObj> createMp4() {
|
||||
Order order = new Order();
|
||||
order.setComplete("1");
|
||||
ExampleMatcher matcher = ExampleMatcher.matching()
|
||||
.withIgnoreNullValues().withIgnorePaths("createAt", "last", "page", "pageSize");
|
||||
Example<Order> ep = Example.of(order, matcher);
|
||||
Sort s = Sort.by("createAt");
|
||||
Flux<Order> flux = orderRepository.findAll(ep, s);
|
||||
flux.filter(fo -> fo.getFileId() == null).subscribe(fo -> {
|
||||
Mono<List<List<String>>> es = recordRepository.findAllByTraceIdOrderByCreateAtAsc(fo.getTraceId()).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, array.toJSONString());
|
||||
} catch (IOException ex) {
|
||||
ex.printStackTrace();
|
||||
}
|
||||
});
|
||||
});
|
||||
return Mono.just(ResObj.ok());
|
||||
}
|
||||
|
||||
void createVideo(Order fo, String jsonData) throws IOException {
|
||||
|
||||
// 创建临时目录
|
||||
Path tempDir = null;
|
||||
try {
|
||||
tempDir = Files.createTempDirectory("ts");
|
||||
} catch (IOException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
System.out.println("Created temporary directory: " + tempDir);
|
||||
|
||||
try {
|
||||
File tempFile = new File(tempDir.toFile(), "traces.json");
|
||||
FileWriter fw = new FileWriter(tempFile);
|
||||
BufferedWriter bw = new BufferedWriter(fw);
|
||||
bw.write(jsonData);
|
||||
bw.flush();
|
||||
bw.close();
|
||||
fw.close();
|
||||
String[] command = {"ts-node", "/opt/trace-transform/src/index.ts"};
|
||||
ProcessBuilder processBuilder = new ProcessBuilder(command);
|
||||
processBuilder.directory(tempDir.toFile());
|
||||
// 在临时目录下执行命令
|
||||
Process process = processBuilder.start();
|
||||
|
||||
Thread processThread = new Thread(() -> {
|
||||
try {
|
||||
// 读取命令输出
|
||||
BufferedReader reader = new BufferedReader(new InputStreamReader(process.getInputStream()));
|
||||
String line;
|
||||
while ((line = reader.readLine()) != null) {
|
||||
System.out.println(line);
|
||||
}
|
||||
reader.close();
|
||||
// 等待命令执行完成
|
||||
int exitCode = process.waitFor();
|
||||
|
||||
System.out.println("Command exited with code " + exitCode);
|
||||
} catch (InterruptedException | IOException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
});
|
||||
processThread.start();
|
||||
|
||||
// Continue with other tasks here...
|
||||
|
||||
try {
|
||||
processThread.join();
|
||||
} catch (InterruptedException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
|
||||
File file = new File(tempDir.toFile(), "video.mp4");
|
||||
if (file.exists()) {
|
||||
log.info(String.valueOf(file.length()));
|
||||
//byte[] fileBytes = Files.readAllBytes(file.toPath());
|
||||
byte[] fileBytes = FileUtil.readBytes(file);
|
||||
String titles = fo.getOrderId() + ".mp4";
|
||||
UploadFileDTO uploadFileDTO = new UploadFileDTO(fileBytes, "RECALL", 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 (Exception e) {
|
||||
e.printStackTrace();
|
||||
log.error(e.getMessage());
|
||||
} finally {
|
||||
// 删除临时目录
|
||||
// Files.walk(tempDir)
|
||||
// .sorted(Comparator.reverseOrder())
|
||||
// .map(Path::toFile)
|
||||
// .forEach(File::delete);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package com.wabestway.recall.web;
|
||||
|
||||
import com.alibaba.fastjson.JSON;
|
||||
import com.wabestway.recall.model.Order;
|
||||
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.web.bind.annotation.*;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.UUID;
|
||||
|
||||
@Slf4j
|
||||
@RestController
|
||||
@RequestMapping("/track")
|
||||
public class RecallController {
|
||||
private final RecordRepository recordRepository;
|
||||
private final OrderRepository orderRepository;
|
||||
|
||||
public RecallController(RecordRepository recordRepository, OrderRepository orderRepository) {
|
||||
this.recordRepository = recordRepository;
|
||||
this.orderRepository = orderRepository;
|
||||
}
|
||||
|
||||
@GetMapping
|
||||
public Flux<Record> getAllUsers() {
|
||||
return recordRepository.findAll();
|
||||
}
|
||||
|
||||
@PostMapping("/save")
|
||||
public Mono<ResObj> save(@RequestBody Record record) {
|
||||
if (record.getOrderId() != null && !record.getOrderId().isEmpty()) {
|
||||
Mono<Order> order = orderRepository.findByOrderId(record.getOrderId());
|
||||
order.hasElement().subscribe(exit -> {
|
||||
if (!exit) {
|
||||
if (record.getTraceId() != null && !record.getTraceId().isEmpty()) {
|
||||
Mono<Order> tcOrder = orderRepository.findByTraceId(record.getTraceId());
|
||||
tcOrder.filter(to -> to.getId() != null).flatMap(to -> {
|
||||
to.setOrderId(record.getOrderId());
|
||||
if (record.isLast()) {
|
||||
to.setComplete("1");
|
||||
}
|
||||
return orderRepository.save(to);
|
||||
}).subscribe();
|
||||
} else {
|
||||
record.setTraceId(UUID.randomUUID().toString().replace("-", ""));
|
||||
Order forder = new Order();
|
||||
forder.setOrderId(record.getOrderId());
|
||||
forder.setTraceId(record.getTraceId());
|
||||
forder.setProductCode(record.getProductCode());
|
||||
forder.setProductName(record.getProductName());
|
||||
forder.setCreateAt(System.currentTimeMillis());
|
||||
forder.setAppKey(record.getAppKey());
|
||||
forder.setArchived("0");
|
||||
forder.setComplete("0");
|
||||
// forder.setHolderName("王五");
|
||||
// forder.setHolderPhone("18515064530");
|
||||
// forder.setSupplierName("阳光保险");
|
||||
// forder.setStartDate("2023-01-01");
|
||||
// forder.setEndDate("2035-01-01");
|
||||
if (record.isLast()) {
|
||||
forder.setComplete("1");
|
||||
}
|
||||
orderRepository.save(forder).subscribe();
|
||||
}
|
||||
} else {
|
||||
order.subscribe(to -> {
|
||||
if (record.isLast()) {
|
||||
to.setComplete("1");
|
||||
orderRepository.save(to).subscribe();
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
} else {
|
||||
if (record.getTraceId() == null || record.getTraceId().isEmpty()) {
|
||||
record.setTraceId(UUID.randomUUID().toString().replace("-", ""));
|
||||
Order order = new Order();
|
||||
order.setTraceId(record.getTraceId());
|
||||
order.setProductCode(record.getProductCode());
|
||||
order.setProductName(record.getProductName());
|
||||
order.setCreateAt(System.currentTimeMillis());
|
||||
order.setAppKey(record.getAppKey());
|
||||
order.setArchived("0");
|
||||
order.setComplete("0");
|
||||
// order.setHolderName("王五");
|
||||
// order.setHolderPhone("18515064530");
|
||||
// order.setSupplierName("阳光保险");
|
||||
// order.setStartDate("2023-01-01");
|
||||
// order.setEndDate("2035-01-01");
|
||||
if (record.isLast()) {
|
||||
order.setComplete("1");
|
||||
}
|
||||
orderRepository.save(order).subscribe();
|
||||
}
|
||||
}
|
||||
record.setCreateAt(System.currentTimeMillis());
|
||||
log.info("保存数据:{}", JSON.toJSONString(record));
|
||||
return recordRepository.save(record).flatMap(ss -> {
|
||||
ss.setEvents(null);
|
||||
return Mono.just(ResObj.ok(ss));
|
||||
});
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user