修复金综流水多文件拉取处理
This commit is contained in:
@@ -686,13 +686,28 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
throw new RuntimeException("拉取本行信息失败: 未返回logId");
|
||||
}
|
||||
|
||||
Integer logId = response.getData().get(0);
|
||||
if (logId == null) {
|
||||
List<Integer> logIds = response.getData().stream()
|
||||
.filter(Objects::nonNull)
|
||||
.toList();
|
||||
if (logIds.isEmpty()) {
|
||||
throw new RuntimeException("拉取本行信息失败: 未返回logId");
|
||||
}
|
||||
|
||||
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, caller);
|
||||
return true;
|
||||
boolean anySuccess = false;
|
||||
for (int i = 0; i < logIds.size(); i++) {
|
||||
Integer logId = logIds.get(i);
|
||||
CcdiFileUploadRecord currentRecord = i == 0
|
||||
? record
|
||||
: createAdditionalPullBankInfoRecord(record, idCard);
|
||||
try {
|
||||
anySuccess |= processRecordAfterLogIdReady(projectId, lsfxProjectId, currentRecord, logId, caller);
|
||||
} catch (Exception logException) {
|
||||
log.error("【拉取本行信息】处理logId失败: idCard={}, logId={}, recordId={}",
|
||||
idCard, logId, currentRecord.getId(), logException);
|
||||
updateFailedRecord(currentRecord, logException.getMessage());
|
||||
}
|
||||
}
|
||||
return anySuccess;
|
||||
} catch (Exception e) {
|
||||
log.error("【拉取本行信息】处理失败: idCard={}, recordId={}", idCard, record.getId(), e);
|
||||
updateFailedRecord(record, e.getMessage());
|
||||
@@ -700,6 +715,24 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
}
|
||||
}
|
||||
|
||||
private CcdiFileUploadRecord createAdditionalPullBankInfoRecord(CcdiFileUploadRecord sourceRecord,
|
||||
String idCard) {
|
||||
CcdiFileUploadRecord record = new CcdiFileUploadRecord();
|
||||
record.setProjectId(sourceRecord.getProjectId());
|
||||
record.setLsfxProjectId(sourceRecord.getLsfxProjectId());
|
||||
record.setFileName(idCard);
|
||||
record.setFileSize(0L);
|
||||
record.setFileStatus("uploading");
|
||||
record.setAccountNos(idCard);
|
||||
record.setUploadTime(new Date());
|
||||
record.setUploadUser(sourceRecord.getUploadUser());
|
||||
recordMapper.insertBatch(List.of(record));
|
||||
if (record.getId() == null) {
|
||||
throw new RuntimeException("创建金综流水上传记录失败: 未生成记录ID");
|
||||
}
|
||||
return record;
|
||||
}
|
||||
|
||||
/**
|
||||
* 异步处理单个文件的完整流程
|
||||
* 包含:上传 → 轮询解析状态 → 获取结果 → 保存流水数据
|
||||
@@ -782,20 +815,20 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
bankTagService.submitAutoRebuild(projectId, triggerType);
|
||||
}
|
||||
|
||||
private void processRecordAfterLogIdReady(Long projectId,
|
||||
Integer lsfxProjectId,
|
||||
CcdiFileUploadRecord record,
|
||||
Integer logId,
|
||||
CallerContext caller) {
|
||||
processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, false, caller);
|
||||
private boolean processRecordAfterLogIdReady(Long projectId,
|
||||
Integer lsfxProjectId,
|
||||
CcdiFileUploadRecord record,
|
||||
Integer logId,
|
||||
CallerContext caller) {
|
||||
return processRecordAfterLogIdReady(projectId, lsfxProjectId, record, logId, false, caller);
|
||||
}
|
||||
|
||||
private void processRecordAfterLogIdReady(Long projectId,
|
||||
Integer lsfxProjectId,
|
||||
CcdiFileUploadRecord record,
|
||||
Integer logId,
|
||||
boolean preserveRecordFileName,
|
||||
CallerContext caller) {
|
||||
private boolean processRecordAfterLogIdReady(Long projectId,
|
||||
Integer lsfxProjectId,
|
||||
CcdiFileUploadRecord record,
|
||||
Integer logId,
|
||||
boolean preserveRecordFileName,
|
||||
CallerContext caller) {
|
||||
log.info("【文件上传】步骤3: 更新状态为解析中, logId={}", logId);
|
||||
record.setLogId(logId);
|
||||
record.setFileStatus("parsing");
|
||||
@@ -840,7 +873,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
if (!parseSuccess) {
|
||||
log.warn("【文件上传】步骤6: 解析失败: status={}, desc={}", status, uploadStatusDesc);
|
||||
updateFailedRecord(record, "解析失败: " + uploadStatusDesc);
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
|
||||
log.info("【文件上传】步骤6: 解析成功,保存主体信息");
|
||||
@@ -857,7 +890,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
logId, fallbackCretNo);
|
||||
if (!fetchResult.isSuccess()) {
|
||||
updateFailedRecord(record, fetchResult.getErrorMessage());
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
|
||||
record.setFileStatus("parsed_success");
|
||||
@@ -865,6 +898,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
|
||||
record.setAccountNos(accountNosStr);
|
||||
record.setErrorMessage(null);
|
||||
recordMapper.updateById(record);
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -510,6 +510,65 @@ class CcdiFileUploadServiceImplTest {
|
||||
));
|
||||
}
|
||||
|
||||
@Test
|
||||
void processPullBankInfoAsync_shouldProcessAllJzlLogIds() {
|
||||
Integer secondLogId = LOG_ID + 1;
|
||||
Integer thirdLogId = LOG_ID + 2;
|
||||
List<CcdiFileUploadRecord> insertedAdditionalRecords = new ArrayList<>();
|
||||
|
||||
doAnswer(invocation -> {
|
||||
List<CcdiFileUploadRecord> records = invocation.getArgument(0);
|
||||
for (int i = 0; i < records.size(); i++) {
|
||||
records.get(i).setId(RECORD_ID + i + 1);
|
||||
CcdiFileUploadRecord snapshot = new CcdiFileUploadRecord();
|
||||
snapshot.setId(records.get(i).getId());
|
||||
snapshot.setFileName(records.get(i).getFileName());
|
||||
snapshot.setAccountNos(records.get(i).getAccountNos());
|
||||
snapshot.setUploadUser(records.get(i).getUploadUser());
|
||||
snapshot.setFileStatus(records.get(i).getFileStatus());
|
||||
insertedAdditionalRecords.add(snapshot);
|
||||
}
|
||||
return records.size();
|
||||
}).when(recordMapper).insertBatch(any());
|
||||
|
||||
when(lsfxClient.fetchInnerFlow(eq(CALLER), any()))
|
||||
.thenReturn(buildFetchInnerFlowResponse(LOG_ID, secondLogId, thirdLogId));
|
||||
when(lsfxClient.checkParseStatus(eq(CALLER), eq(LSFX_PROJECT_ID), org.mockito.ArgumentMatchers.anyString()))
|
||||
.thenReturn(buildCheckParseStatusResponse(false));
|
||||
when(lsfxClient.getFileUploadStatus(eq(CALLER), any())).thenReturn(buildParsedSuccessStatusResponse());
|
||||
when(lsfxClient.getBankStatement(eq(CALLER), any(GetBankStatementRequest.class)))
|
||||
.thenReturn(buildEmptyBankStatementResponse());
|
||||
|
||||
CcdiFileUploadRecord record = buildRecord();
|
||||
record.setUploadUser("admin");
|
||||
|
||||
boolean success = service.processPullBankInfoAsync(
|
||||
PROJECT_ID,
|
||||
LSFX_PROJECT_ID,
|
||||
record,
|
||||
"110101199001018888",
|
||||
LsfxConstants.DATA_CHANNEL_JZL,
|
||||
null,
|
||||
null,
|
||||
CALLER
|
||||
);
|
||||
|
||||
assertTrue(success);
|
||||
assertEquals(2, insertedAdditionalRecords.size());
|
||||
assertEquals("110101199001018888", insertedAdditionalRecords.get(0).getFileName());
|
||||
assertEquals("110101199001018888", insertedAdditionalRecords.get(0).getAccountNos());
|
||||
assertEquals("admin", insertedAdditionalRecords.get(0).getUploadUser());
|
||||
verify(lsfxClient).checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(LOG_ID));
|
||||
verify(lsfxClient).checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(secondLogId));
|
||||
verify(lsfxClient).checkParseStatus(CALLER, LSFX_PROJECT_ID, String.valueOf(thirdLogId));
|
||||
verify(lsfxClient).getBankStatement(eq(CALLER), org.mockito.ArgumentMatchers.<GetBankStatementRequest>argThat(request ->
|
||||
LOG_ID.equals(request.getLogId())));
|
||||
verify(lsfxClient).getBankStatement(eq(CALLER), org.mockito.ArgumentMatchers.<GetBankStatementRequest>argThat(request ->
|
||||
secondLogId.equals(request.getLogId())));
|
||||
verify(lsfxClient).getBankStatement(eq(CALLER), org.mockito.ArgumentMatchers.<GetBankStatementRequest>argThat(request ->
|
||||
thirdLogId.equals(request.getLogId())));
|
||||
}
|
||||
|
||||
@Test
|
||||
void processFileAsync_shouldUploadToLsfxWithOriginalRecordFileName() throws IOException {
|
||||
when(lsfxClient.uploadFile(eq(CALLER), eq(LSFX_PROJECT_ID), any(), eq("原始流水.xlsx")))
|
||||
@@ -995,9 +1054,9 @@ class CcdiFileUploadServiceImplTest {
|
||||
return response;
|
||||
}
|
||||
|
||||
private FetchInnerFlowResponse buildFetchInnerFlowResponse(Integer logId) {
|
||||
private FetchInnerFlowResponse buildFetchInnerFlowResponse(Integer... logIds) {
|
||||
FetchInnerFlowResponse response = new FetchInnerFlowResponse();
|
||||
response.setData(List.of(logId));
|
||||
response.setData(List.of(logIds));
|
||||
return response;
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user