feat: 接入新华社工商信息同步

This commit is contained in:
wkc
2026-07-29 17:56:27 +08:00
parent 2bcba71259
commit 5d004a66e8
72 changed files with 2610 additions and 59 deletions

View File

@@ -30,6 +30,13 @@
<version>${ruoyi.version}</version>
</dependency>
<!-- 企业工商信息统一查询能力 -->
<dependency>
<groupId>com.ruoyi</groupId>
<artifactId>ccdi-info-collection</artifactId>
<version>${ruoyi.version}</version>
</dependency>
<!-- lombok -->
<dependency>
<groupId>org.projectlombok</groupId>

View File

@@ -186,8 +186,8 @@ public class CcdiFileUploadController extends BaseController {
@PreAuthorize("@ss.hasPermi('ccdi:project:edit')")
public AjaxResult deleteFile(@PathVariable Long id) {
projectAccessService.assertCanOperateByFileRecordId(id);
Long userId = SecurityUtils.getUserId();
String message = fileUploadService.deleteFileUploadRecord(id, userId);
CallerContext caller = CallerContext.from(SecurityUtils.getLoginUser());
String message = fileUploadService.deleteFileUploadRecord(id, caller);
return AjaxResult.success(message);
}
}

View File

@@ -0,0 +1,56 @@
package com.ruoyi.ccdi.project.controller;
import com.ruoyi.ccdi.project.domain.dto.ProjectCounterpartyRefreshDTO;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.common.core.domain.AjaxResult;
import com.ruoyi.common.utils.SecurityUtils;
import com.ruoyi.lsfx.domain.CallerContext;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.tags.Tag;
import jakarta.annotation.Resource;
import jakarta.validation.Valid;
import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
@RestController
@RequestMapping("/ccdi/project/counterparty-enterprise")
@Tag(name = "项目对手方工商信息")
public class CcdiProjectCounterpartyEnterpriseController {
@Resource
private IProjectCounterpartyEnterpriseService service;
@GetMapping("/detail")
@Operation(summary = "查询本地项目对手方工商详情")
@PreAuthorize("@ss.hasPermi('ccdi:project:query')")
public AjaxResult detail(@RequestParam Long projectId, @RequestParam String counterpartyName) {
return AjaxResult.success(service.getDetail(projectId, counterpartyName));
}
@PostMapping("/refresh")
@Operation(summary = "重新查询单个项目对手方工商信息")
@PreAuthorize("@ss.hasPermi('ccdi:project:counterpartyEnterprise:refresh')")
public AjaxResult refresh(@Valid @RequestBody ProjectCounterpartyRefreshDTO dto) {
return AjaxResult.success(service.refresh(CallerContext.from(SecurityUtils.getLoginUser()),
dto.getProjectId(), dto.getCounterpartyName()));
}
@PostMapping("/backfill")
@Operation(summary = "启动项目对手方工商信息历史补全")
@PreAuthorize("@ss.hasPermi('ccdi:project:counterpartyEnterprise:backfill')")
public AjaxResult backfill() {
return AjaxResult.success("任务已提交", service.startBackfill(CallerContext.from(SecurityUtils.getLoginUser())));
}
@GetMapping("/backfill/{taskId}")
@Operation(summary = "查询项目对手方工商信息补全状态")
@PreAuthorize("@ss.hasPermi('ccdi:project:counterpartyEnterprise:backfill')")
public AjaxResult backfillStatus(@PathVariable String taskId) {
return AjaxResult.success(service.getBackfillStatus(taskId));
}
}

View File

@@ -0,0 +1,14 @@
package com.ruoyi.ccdi.project.domain.dto;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull;
import lombok.Data;
@Data
public class ProjectCounterpartyRefreshDTO {
@NotNull(message = "项目ID不能为空")
private Long projectId;
@NotBlank(message = "对手方名称不能为空")
private String counterpartyName;
}

View File

@@ -0,0 +1,41 @@
package com.ruoyi.ccdi.project.domain.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.Date;
/** 项目对手方工商信息。 */
@Data
@TableName("ccdi_project_counterparty_enterprise")
public class CcdiProjectCounterpartyEnterprise {
@TableId(type = IdType.AUTO)
private Long counterpartyEnterpriseId;
private Long projectId;
private String counterpartyName;
private String socialCreditCode;
private String enterpriseName;
private BigDecimal registeredCapital;
private String registeredCapitalUnit;
private LocalDate registerDate;
private LocalDate establishDate;
private String industryCode;
private String industryName;
private String organizationTypeCode;
private String organizationTypeName;
private String regionCode;
private String regionName;
private String registerAddress;
private Integer employeeCount;
private String legalRepresentative;
private Long cacheInfoId;
private Date syncTime;
private String createBy;
private Date createTime;
private String updateBy;
private Date updateTime;
}

View File

@@ -0,0 +1,26 @@
package com.ruoyi.ccdi.project.domain.entity;
import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableField;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import lombok.Data;
import java.math.BigDecimal;
import java.util.Date;
/** 项目对手方股东信息。 */
@Data
@TableName("ccdi_project_counterparty_shareholder")
public class CcdiProjectCounterpartyShareholder {
@TableId(type = IdType.AUTO)
private Long shareholderId;
private Long counterpartyEnterpriseId;
@TableField("shareholder_seq")
private Integer sequenceNo;
private String shareholderName;
private BigDecimal stockPercent;
private BigDecimal subscribedCapital;
private String capitalUnit;
private Date createTime;
}

View File

@@ -1,6 +1,7 @@
package com.ruoyi.ccdi.project.domain.event;
import com.ruoyi.ccdi.project.domain.dto.CcdiProjectImportHistoryDTO;
import com.ruoyi.lsfx.domain.CallerContext;
/**
* 历史项目导入提交事件
@@ -10,14 +11,14 @@ public class CcdiProjectHistoryImportSubmittedEvent {
private final Long targetProjectId;
private final Integer targetLsfxProjectId;
private final CcdiProjectImportHistoryDTO dto;
private final String operator;
private final CallerContext caller;
public CcdiProjectHistoryImportSubmittedEvent(Long targetProjectId, Integer targetLsfxProjectId,
CcdiProjectImportHistoryDTO dto, String operator) {
CcdiProjectImportHistoryDTO dto, CallerContext caller) {
this.targetProjectId = targetProjectId;
this.targetLsfxProjectId = targetLsfxProjectId;
this.dto = dto;
this.operator = operator;
this.caller = caller;
}
public Long getTargetProjectId() {
@@ -32,7 +33,7 @@ public class CcdiProjectHistoryImportSubmittedEvent {
return dto;
}
public String getOperator() {
return operator;
public CallerContext getCaller() {
return caller;
}
}

View File

@@ -0,0 +1,35 @@
package com.ruoyi.ccdi.project.domain.vo;
import lombok.Data;
import java.math.BigDecimal;
import java.time.LocalDate;
import java.util.ArrayList;
import java.util.Date;
import java.util.List;
@Data
public class ProjectCounterpartyEnterpriseVO {
private Long projectId;
private String counterpartyName;
private String status;
private Boolean canRefresh;
private Date cacheValidDate;
private String socialCreditCode;
private String enterpriseName;
private BigDecimal registeredCapital;
private String registeredCapitalUnit;
private LocalDate registerDate;
private LocalDate establishDate;
private String industryCode;
private String industryName;
private String organizationTypeCode;
private String organizationTypeName;
private String regionCode;
private String regionName;
private String registerAddress;
private Integer employeeCount;
private String legalRepresentative;
private Date syncTime;
private List<ProjectCounterpartyShareholderVO> shareholders = new ArrayList<>();
}

View File

@@ -0,0 +1,14 @@
package com.ruoyi.ccdi.project.domain.vo;
import lombok.Data;
import java.math.BigDecimal;
@Data
public class ProjectCounterpartyShareholderVO {
private Integer sequenceNo;
private String shareholderName;
private BigDecimal stockPercent;
private BigDecimal subscribedCapital;
private String capitalUnit;
}

View File

@@ -0,0 +1,15 @@
package com.ruoyi.ccdi.project.domain.vo;
import lombok.Data;
@Data
public class ProjectCounterpartySyncTaskVO {
private String taskId;
private String status;
private Integer totalCount;
private Integer completedCount;
private Integer successCount;
private Integer skippedCount;
private Integer failureCount;
private String message;
}

View File

@@ -49,4 +49,8 @@ public interface CcdiBankStatementMapper extends BaseMapper<CcdiBankStatement> {
CcdiBankStatementFilterOptionsVO selectFilterOptions(@Param("projectId") Long projectId);
Integer countMatchedStaffCountByProjectId(@Param("projectId") Long projectId);
List<String> selectDistinctCounterpartyNames(@Param("projectId") Long projectId);
List<Long> selectDistinctProjectIdsWithCounterparties();
}

View File

@@ -0,0 +1,16 @@
package com.ruoyi.ccdi.project.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyEnterprise;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface CcdiProjectCounterpartyEnterpriseMapper extends BaseMapper<CcdiProjectCounterpartyEnterprise> {
CcdiProjectCounterpartyEnterprise selectByProjectAndName(@Param("projectId") Long projectId,
@Param("counterpartyName") String counterpartyName);
List<CcdiProjectCounterpartyEnterprise> selectByProjectId(@Param("projectId") Long projectId);
int deleteByProjectId(@Param("projectId") Long projectId);
}

View File

@@ -0,0 +1,15 @@
package com.ruoyi.ccdi.project.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyShareholder;
import org.apache.ibatis.annotations.Param;
import java.util.List;
public interface CcdiProjectCounterpartyShareholderMapper extends BaseMapper<CcdiProjectCounterpartyShareholder> {
List<CcdiProjectCounterpartyShareholder> selectByEnterpriseId(@Param("counterpartyEnterpriseId") Long counterpartyEnterpriseId);
int deleteByEnterpriseId(@Param("counterpartyEnterpriseId") Long counterpartyEnterpriseId);
int deleteByProjectId(@Param("projectId") Long projectId);
}

View File

@@ -58,10 +58,10 @@ public interface ICcdiFileUploadService {
* 删除上传记录并清理关联数据
*
* @param id 上传记录ID
* @param operatorUserId 当前操作用户ID
* @param caller 当前操作人上下文
* @return 删除结果
*/
String deleteFileUploadRecord(Long id, Long operatorUserId);
String deleteFileUploadRecord(Long id, CallerContext caller);
/**
* 查询上传记录列表

View File

@@ -1,6 +1,7 @@
package com.ruoyi.ccdi.project.service;
import com.ruoyi.ccdi.project.domain.dto.CcdiProjectImportHistoryDTO;
import com.ruoyi.lsfx.domain.CallerContext;
/**
* 历史项目导入服务
@@ -15,8 +16,8 @@ public interface ICcdiProjectHistoryImportService {
* @param targetProjectId 目标项目ID
* @param targetLsfxProjectId 目标流水分析项目ID
* @param dto 导入参数
* @param operator 操作人
* @param caller 原始操作人上下文
*/
void submitImport(Long targetProjectId, Integer targetLsfxProjectId,
CcdiProjectImportHistoryDTO dto, String operator);
CcdiProjectImportHistoryDTO dto, CallerContext caller);
}

View File

@@ -0,0 +1,14 @@
package com.ruoyi.ccdi.project.service;
import com.ruoyi.ccdi.project.domain.vo.ProjectCounterpartyEnterpriseVO;
import com.ruoyi.ccdi.project.domain.vo.ProjectCounterpartySyncTaskVO;
import com.ruoyi.lsfx.domain.CallerContext;
public interface IProjectCounterpartyEnterpriseService {
ProjectCounterpartyEnterpriseVO getDetail(Long projectId, String counterpartyName);
ProjectCounterpartyEnterpriseVO refresh(CallerContext caller, Long projectId, String counterpartyName);
void submitReconcile(CallerContext caller, Long projectId);
String startBackfill(CallerContext caller);
ProjectCounterpartySyncTaskVO getBackfillStatus(String taskId);
void deleteProjectData(Long projectId);
}

View File

@@ -16,6 +16,7 @@ import com.ruoyi.ccdi.project.mapper.CcdiProjectMapper;
import com.ruoyi.ccdi.project.service.ICcdiBankTagService;
import com.ruoyi.ccdi.project.service.ICcdiFileUploadService;
import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.constants.LsfxConstants;
@@ -106,6 +107,9 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
@Resource
private ICcdiProjectService projectService;
@Resource
private IProjectCounterpartyEnterpriseService counterpartyEnterpriseService;
/**
* 获取临时文件存储目录
*/
@@ -243,7 +247,8 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
}
@Override
public String deleteFileUploadRecord(Long id, Long operatorUserId) {
@Transactional(rollbackFor = Exception.class)
public String deleteFileUploadRecord(Long id, CallerContext caller) {
CcdiFileUploadRecord record = recordMapper.selectById(id);
validateDeleteRecord(record);
@@ -252,7 +257,7 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
* DeleteFilesRequest request = new DeleteFilesRequest();
* request.setGroupId(record.getLsfxProjectId());
* request.setLogIds(new Integer[]{record.getLogId()});
* request.setUserId(toUploadUserId(operatorUserId));
* request.setUserId(toUploadUserId(caller.userId()));
*
* DeleteFilesResponse response = lsfxClient.deleteFiles(request);
* if (response == null || Boolean.FALSE.equals(response.getSuccessResponse())) {
@@ -272,9 +277,23 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
}
bankTagService.submitAutoRebuild(record.getProjectId(), TriggerType.AUTO_FILE_DELETE);
submitCounterpartyReconcileAfterCommit(caller, record.getProjectId());
return "删除成功,已开始项目重新打标";
}
private void submitCounterpartyReconcileAfterCommit(CallerContext caller, Long projectId) {
if (!TransactionSynchronizationManager.isSynchronizationActive()) {
counterpartyEnterpriseService.submitReconcile(caller, projectId);
return;
}
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
counterpartyEnterpriseService.submitReconcile(caller, projectId);
}
});
}
@Override
public Page<CcdiFileUploadRecord> selectPage(Page<CcdiFileUploadRecord> page,
CcdiFileUploadQueryDTO queryDTO) {
@@ -554,6 +573,9 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
.whenComplete((unused, throwable) -> {
boolean anySuccess = futures.stream().anyMatch(future -> Boolean.TRUE.equals(future.getNow(Boolean.FALSE)));
handleTagRebuildAfterBatchCompletion(projectId, TriggerType.AUTO_BATCH_UPLOAD, anySuccess);
if (anySuccess) {
counterpartyEnterpriseService.submitReconcile(caller, projectId);
}
});
log.info("【文件上传】调度线程完成: projectId={}, batchId={}", projectId, batchId);
@@ -654,6 +676,9 @@ public class CcdiFileUploadServiceImpl implements ICcdiFileUploadService {
.whenComplete((unused, throwable) -> {
boolean anySuccess = futures.stream().anyMatch(future -> Boolean.TRUE.equals(future.getNow(Boolean.FALSE)));
handleTagRebuildAfterBatchCompletion(projectId, TriggerType.AUTO_PULL_BANK_INFO, anySuccess);
if (anySuccess) {
counterpartyEnterpriseService.submitReconcile(caller, projectId);
}
});
}

View File

@@ -21,7 +21,7 @@ public class CcdiProjectHistoryImportEventListener {
event.getTargetProjectId(),
event.getTargetLsfxProjectId(),
event.getDto(),
event.getOperator()
event.getCaller()
);
}
}

View File

@@ -10,6 +10,8 @@ import com.ruoyi.ccdi.project.mapper.CcdiFileUploadRecordMapper;
import com.ruoyi.ccdi.project.mapper.CcdiProjectMapper;
import com.ruoyi.ccdi.project.service.ICcdiBankTagService;
import com.ruoyi.ccdi.project.service.ICcdiProjectHistoryImportService;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.lsfx.domain.CallerContext;
import jakarta.annotation.Resource;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.BeanUtils;
@@ -52,14 +54,17 @@ public class CcdiProjectHistoryImportServiceImpl implements ICcdiProjectHistoryI
@Resource
private ICcdiBankTagService bankTagService;
@Resource
private IProjectCounterpartyEnterpriseService counterpartyEnterpriseService;
@Override
public void submitImport(Long targetProjectId, Integer targetLsfxProjectId,
CcdiProjectImportHistoryDTO dto, String operator) {
fileUploadExecutor.execute(() -> executeImport(targetProjectId, targetLsfxProjectId, dto, operator));
CcdiProjectImportHistoryDTO dto, CallerContext caller) {
fileUploadExecutor.execute(() -> executeImport(targetProjectId, targetLsfxProjectId, dto, caller));
}
private void executeImport(Long targetProjectId, Integer targetLsfxProjectId,
CcdiProjectImportHistoryDTO dto, String operator) {
CcdiProjectImportHistoryDTO dto, CallerContext caller) {
List<CcdiFileUploadRecord> sourceRecords = recordMapper.selectSuccessfulRecordsByProjectIds(dto.getSourceProjectIds());
if (sourceRecords == null || sourceRecords.isEmpty()) {
log.info("【项目历史导入】无可复制的来源批次: projectId={}, sourceProjectIds={}",
@@ -91,7 +96,7 @@ public class CcdiProjectHistoryImportServiceImpl implements ICcdiProjectHistoryI
if (statementsToInsert.size() > sizeBefore) {
recordsToInsert.add(buildHistoryImportRecord(
sourceRecord, targetProjectId, targetLsfxProjectId, newBatchId, resolveSourceProjectName(sourceRecord.getProjectId()), operator
sourceRecord, targetProjectId, targetLsfxProjectId, newBatchId, resolveSourceProjectName(sourceRecord.getProjectId()), caller.username()
));
}
}
@@ -105,6 +110,7 @@ public class CcdiProjectHistoryImportServiceImpl implements ICcdiProjectHistoryI
if (!statementsToInsert.isEmpty()) {
refreshProjectTargetCount(targetProjectId);
bankTagService.submitAutoRebuild(targetProjectId, TriggerType.AUTO_BATCH_UPLOAD);
counterpartyEnterpriseService.submitReconcile(caller, targetProjectId);
}
}

View File

@@ -175,7 +175,7 @@ public class CcdiProjectServiceImpl implements ICcdiProjectService {
@Override
public void afterCommit() {
applicationEventPublisher.publishEvent(
new CcdiProjectHistoryImportSubmittedEvent(project.getProjectId(), project.getLsfxProjectId(), dto, caller.username())
new CcdiProjectHistoryImportSubmittedEvent(project.getProjectId(), project.getLsfxProjectId(), dto, caller)
);
}
});

View File

@@ -0,0 +1,280 @@
package com.ruoyi.ccdi.project.service.impl;
import com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyEnterprise;
import com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyShareholder;
import com.ruoyi.ccdi.project.domain.vo.ProjectCounterpartyEnterpriseVO;
import com.ruoyi.ccdi.project.domain.vo.ProjectCounterpartyShareholderVO;
import com.ruoyi.ccdi.project.domain.vo.ProjectCounterpartySyncTaskVO;
import com.ruoyi.ccdi.project.mapper.CcdiBankStatementMapper;
import com.ruoyi.ccdi.project.mapper.CcdiProjectCounterpartyEnterpriseMapper;
import com.ruoyi.ccdi.project.mapper.CcdiProjectCounterpartyShareholderMapper;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.info.collection.domain.CcdiEnterpriseInfoQueryCache;
import com.ruoyi.info.collection.domain.model.EnterpriseProfile;
import com.ruoyi.info.collection.domain.model.EnterpriseProfileQueryResult;
import com.ruoyi.info.collection.domain.model.EnterpriseShareholderProfile;
import com.ruoyi.info.collection.service.IEnterpriseProfileQueryService;
import com.ruoyi.lsfx.domain.CallerContext;
import jakarta.annotation.Resource;
import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;
import org.springframework.transaction.support.TransactionTemplate;
import org.springframework.util.StringUtils;
import java.util.ArrayList;
import java.util.Date;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@Service
public class ProjectCounterpartyEnterpriseServiceImpl implements IProjectCounterpartyEnterpriseService {
private static final String TASK_PREFIX = "project:counterparty:profile:backfill:";
@Resource
private CcdiBankStatementMapper bankStatementMapper;
@Resource
private CcdiProjectCounterpartyEnterpriseMapper enterpriseMapper;
@Resource
private CcdiProjectCounterpartyShareholderMapper shareholderMapper;
@Resource
private IEnterpriseProfileQueryService queryService;
@Resource
private TransactionTemplate transactionTemplate;
@Resource
private RedisTemplate<String, Object> redisTemplate;
@Resource
@Qualifier("enterpriseProfileExecutor")
private Executor enterpriseProfileExecutor;
@Override
public ProjectCounterpartyEnterpriseVO getDetail(Long projectId, String counterpartyName) {
String name = normalize(counterpartyName);
requireCurrentCounterparty(projectId, name);
CcdiProjectCounterpartyEnterprise entity = enterpriseMapper.selectByProjectAndName(projectId, name);
CcdiEnterpriseInfoQueryCache cache = queryService.findCache(name);
Date now = new Date();
if (entity == null) {
ProjectCounterpartyEnterpriseVO vo = new ProjectCounterpartyEnterpriseVO();
vo.setProjectId(projectId);
vo.setCounterpartyName(name);
vo.setStatus("NOT_FOUND");
vo.setCanRefresh(true);
vo.setCacheValidDate(cache == null ? null : cache.getValidDate());
return vo;
}
ProjectCounterpartyEnterpriseVO vo = new ProjectCounterpartyEnterpriseVO();
BeanUtils.copyProperties(entity, vo);
vo.setStatus(cache != null && cache.getValidDate() != null && cache.getValidDate().after(now) ? "SUCCESS" : "EXPIRED");
vo.setCanRefresh(!"SUCCESS".equals(vo.getStatus()));
vo.setCacheValidDate(cache == null ? null : cache.getValidDate());
vo.setShareholders(shareholderMapper.selectByEnterpriseId(entity.getCounterpartyEnterpriseId()).stream()
.map(this::toShareholderVO).toList());
return vo;
}
@Override
public ProjectCounterpartyEnterpriseVO refresh(CallerContext caller, Long projectId, String counterpartyName) {
String name = normalize(counterpartyName);
requireCurrentCounterparty(projectId, name);
CcdiProjectCounterpartyEnterprise existing = enterpriseMapper.selectByProjectAndName(projectId, name);
CcdiEnterpriseInfoQueryCache cache = queryService.findCache(name);
boolean cacheValid = cache != null && cache.getValidDate() != null && cache.getValidDate().after(new Date());
if (existing != null && cacheValid) {
throw new IllegalStateException("工商缓存仍在有效期内,无需重新查询");
}
EnterpriseProfileQueryResult result = queryService.query(caller, name, !cacheValid);
upsert(projectId, name, result, caller.username());
return getDetail(projectId, name);
}
@Override
public void submitReconcile(CallerContext caller, Long projectId) {
CompletableFuture.runAsync(() -> reconcile(caller, projectId, null));
}
@Override
public String startBackfill(CallerContext caller) {
String taskId = UUID.randomUUID().toString().replace("-", "");
List<Long> projectIds = bankStatementMapper.selectDistinctProjectIdsWithCounterparties();
int total = projectIds.stream().mapToInt(id -> bankStatementMapper.selectDistinctCounterpartyNames(id).size()).sum();
initializeTask(taskId, total);
CompletableFuture.runAsync(() -> {
TaskCounters counters = new TaskCounters();
for (Long projectId : projectIds) {
reconcile(caller, projectId, new Progress(taskId, counters));
}
updateTask(taskId, "COMPLETED", counters);
});
return taskId;
}
@Override
public ProjectCounterpartySyncTaskVO getBackfillStatus(String taskId) {
Map<Object, Object> values = redisTemplate.opsForHash().entries(taskKey(taskId));
if (values.isEmpty()) {
throw new IllegalArgumentException("任务不存在或已过期");
}
ProjectCounterpartySyncTaskVO vo = new ProjectCounterpartySyncTaskVO();
vo.setTaskId(string(values.get("taskId")));
vo.setStatus(string(values.get("status")));
vo.setTotalCount(number(values.get("totalCount")));
vo.setCompletedCount(number(values.get("completedCount")));
vo.setSuccessCount(number(values.get("successCount")));
vo.setSkippedCount(number(values.get("skippedCount")));
vo.setFailureCount(number(values.get("failureCount")));
vo.setMessage(string(values.get("message")));
return vo;
}
@Override
public void deleteProjectData(Long projectId) {
shareholderMapper.deleteByProjectId(projectId);
enterpriseMapper.deleteByProjectId(projectId);
}
private void reconcile(CallerContext caller, Long projectId, Progress progress) {
List<String> names = bankStatementMapper.selectDistinctCounterpartyNames(projectId);
if (names.isEmpty()) {
transactionTemplate.executeWithoutResult(status -> deleteProjectData(projectId));
return;
}
for (int start = 0; start < names.size(); start += 100) {
List<String> batch = names.subList(start, Math.min(start + 100, names.size()));
List<CompletableFuture<Void>> futures = batch.stream().map(name -> CompletableFuture.runAsync(() -> {
try {
EnterpriseProfileQueryResult result = queryService.query(caller, name, false);
upsert(projectId, name, result, caller.username());
if (progress != null) progress.counters.success.incrementAndGet();
} catch (Exception e) {
if (progress != null) progress.counters.failure.incrementAndGet();
} finally {
if (progress != null) {
progress.counters.completed.incrementAndGet();
updateTask(progress.taskId, "PROCESSING", progress.counters);
}
}
}, enterpriseProfileExecutor)).toList();
CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join();
}
removeStale(projectId, new HashSet<>(names));
}
private void upsert(Long projectId, String name, EnterpriseProfileQueryResult result, String username) {
transactionTemplate.executeWithoutResult(status -> {
EnterpriseProfile profile = result.profile();
CcdiProjectCounterpartyEnterprise entity = enterpriseMapper.selectByProjectAndName(projectId, name);
Date now = new Date();
boolean insert = entity == null;
if (insert) {
entity = new CcdiProjectCounterpartyEnterprise();
entity.setProjectId(projectId);
entity.setCounterpartyName(name);
entity.setCreateBy(username);
entity.setCreateTime(now);
}
entity.setSocialCreditCode(profile.getCreditCode());
entity.setEnterpriseName(profile.getEnterpriseName());
entity.setRegisteredCapital(profile.getRegisteredCapital());
entity.setRegisteredCapitalUnit(profile.getRegisteredCapitalUnit());
entity.setRegisterDate(profile.getRegisterDate());
entity.setEstablishDate(profile.getEstablishDate());
entity.setIndustryCode(profile.getIndustryCode());
entity.setIndustryName(profile.getIndustryName());
entity.setOrganizationTypeCode(profile.getOrganizationTypeCode());
entity.setOrganizationTypeName(profile.getOrganizationTypeName());
entity.setRegionCode(profile.getRegionCode());
entity.setRegionName(profile.getRegionName());
entity.setRegisterAddress(profile.getRegisterAddress());
entity.setEmployeeCount(profile.getEmployeeCount());
entity.setLegalRepresentative(profile.getLegalRepresentative());
entity.setCacheInfoId(result.cacheInfoId());
entity.setSyncTime(now);
entity.setUpdateBy(username);
entity.setUpdateTime(now);
if (insert) enterpriseMapper.insert(entity); else enterpriseMapper.updateById(entity);
shareholderMapper.deleteByEnterpriseId(entity.getCounterpartyEnterpriseId());
for (EnterpriseShareholderProfile item : profile.getShareholders()) {
CcdiProjectCounterpartyShareholder shareholder = new CcdiProjectCounterpartyShareholder();
shareholder.setCounterpartyEnterpriseId(entity.getCounterpartyEnterpriseId());
shareholder.setSequenceNo(item.getSequence());
shareholder.setShareholderName(item.getShareholderName());
shareholder.setStockPercent(item.getStockPercent());
shareholder.setSubscribedCapital(item.getSubscribedCapital());
shareholder.setCapitalUnit(item.getCapitalUnit());
shareholder.setCreateTime(now);
shareholderMapper.insert(shareholder);
}
});
}
private void removeStale(Long projectId, Set<String> currentNames) {
transactionTemplate.executeWithoutResult(status -> enterpriseMapper.selectByProjectId(projectId).stream()
.filter(item -> !currentNames.contains(item.getCounterpartyName()))
.forEach(item -> {
shareholderMapper.deleteByEnterpriseId(item.getCounterpartyEnterpriseId());
enterpriseMapper.deleteById(item.getCounterpartyEnterpriseId());
}));
}
private void requireCurrentCounterparty(Long projectId, String name) {
if (projectId == null || !StringUtils.hasText(name)
|| !bankStatementMapper.selectDistinctCounterpartyNames(projectId).contains(name)) {
throw new IllegalArgumentException("该对手方不在项目当前流水中");
}
}
private ProjectCounterpartyShareholderVO toShareholderVO(CcdiProjectCounterpartyShareholder entity) {
ProjectCounterpartyShareholderVO vo = new ProjectCounterpartyShareholderVO();
BeanUtils.copyProperties(entity, vo);
return vo;
}
private String normalize(String value) { return value == null ? null : value.trim(); }
private String taskKey(String taskId) { return TASK_PREFIX + taskId; }
private String string(Object value) { return value == null ? null : value.toString(); }
private Integer number(Object value) { return value == null ? 0 : Integer.valueOf(value.toString()); }
private void initializeTask(String taskId, int total) {
Map<String, Object> values = new LinkedHashMap<>();
values.put("taskId", taskId);
values.put("status", "PROCESSING");
values.put("totalCount", total);
values.put("completedCount", 0);
values.put("successCount", 0);
values.put("skippedCount", 0);
values.put("failureCount", 0);
values.put("message", "正在补全项目对手方工商信息");
redisTemplate.opsForHash().putAll(taskKey(taskId), values);
redisTemplate.expire(taskKey(taskId), 7, TimeUnit.DAYS);
}
private void updateTask(String taskId, String status, TaskCounters counters) {
Map<String, Object> values = new LinkedHashMap<>();
values.put("status", status);
values.put("completedCount", counters.completed.get());
values.put("successCount", counters.success.get());
values.put("skippedCount", counters.skipped.get());
values.put("failureCount", counters.failure.get());
values.put("message", "COMPLETED".equals(status) ? "项目对手方工商信息补全完成" : "正在补全项目对手方工商信息");
redisTemplate.opsForHash().putAll(taskKey(taskId), values);
}
private static class TaskCounters {
private final AtomicInteger completed = new AtomicInteger();
private final AtomicInteger success = new AtomicInteger();
private final AtomicInteger skipped = new AtomicInteger();
private final AtomicInteger failure = new AtomicInteger();
}
private record Progress(String taskId, TaskCounters counters) { }
}

View File

@@ -139,6 +139,24 @@ PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN"
and trim(bs.cret_no) != ''
</select>
<select id="selectDistinctCounterpartyNames" resultType="java.lang.String">
select distinct trim(CUSTOMER_ACCOUNT_NAME)
from ccdi_bank_statement
where project_id = #{projectId}
and CUSTOMER_ACCOUNT_NAME is not null
and trim(CUSTOMER_ACCOUNT_NAME) != ''
order by trim(CUSTOMER_ACCOUNT_NAME)
</select>
<select id="selectDistinctProjectIdsWithCounterparties" resultType="java.lang.Long">
select distinct project_id
from ccdi_bank_statement
where project_id is not null
and CUSTOMER_ACCOUNT_NAME is not null
and trim(CUSTOMER_ACCOUNT_NAME) != ''
order by project_id
</select>
<sql id="parsedTrxDateExpr">
CASE
WHEN bs.TRX_DATE IS NULL OR TRIM(bs.TRX_DATE) = '' THEN NULL

View File

@@ -0,0 +1,17 @@
<?xml version="1.0" encoding="UTF-8" ?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.ruoyi.ccdi.project.mapper.CcdiProjectCounterpartyEnterpriseMapper">
<select id="selectByProjectAndName" resultType="com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyEnterprise">
select * from ccdi_project_counterparty_enterprise
where project_id = #{projectId} and counterparty_name = #{counterpartyName}
limit 1
</select>
<select id="selectByProjectId" resultType="com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyEnterprise">
select * from ccdi_project_counterparty_enterprise where project_id = #{projectId}
</select>
<delete id="deleteByProjectId">
delete from ccdi_project_counterparty_enterprise where project_id = #{projectId}
</delete>
</mapper>

View File

@@ -0,0 +1,21 @@
<?xml version="1.0" encoding="UTF-8" ?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.ruoyi.ccdi.project.mapper.CcdiProjectCounterpartyShareholderMapper">
<select id="selectByEnterpriseId" resultType="com.ruoyi.ccdi.project.domain.entity.CcdiProjectCounterpartyShareholder">
select * from ccdi_project_counterparty_shareholder
where counterparty_enterprise_id = #{counterpartyEnterpriseId}
order by shareholder_seq, shareholder_id
</select>
<delete id="deleteByEnterpriseId">
delete from ccdi_project_counterparty_shareholder
where counterparty_enterprise_id = #{counterpartyEnterpriseId}
</delete>
<delete id="deleteByProjectId">
delete s from ccdi_project_counterparty_shareholder s
inner join ccdi_project_counterparty_enterprise e
on e.counterparty_enterprise_id = s.counterparty_enterprise_id
where e.project_id = #{projectId}
</delete>
</mapper>

View File

@@ -134,16 +134,16 @@ class CcdiFileUploadControllerTest {
}
@Test
void deleteFile_shouldUseCurrentLoginUserId() {
void deleteFile_shouldUseCurrentLoginUser() {
setLoginUser(9527L, "admin");
when(fileUploadService.deleteFileUploadRecord(123L, 9527L))
when(fileUploadService.deleteFileUploadRecord(123L, CALLER))
.thenReturn("删除成功");
AjaxResult result = controller.deleteFile(123L);
assertEquals(200, result.get("code"));
assertEquals("删除成功", result.get("msg"));
verify(fileUploadService).deleteFileUploadRecord(123L, 9527L);
verify(fileUploadService).deleteFileUploadRecord(123L, CALLER);
}
private void setLoginUser(Long userId, String username) {

View File

@@ -14,6 +14,7 @@ import com.ruoyi.ccdi.project.mapper.CcdiFileUploadRecordMapper;
import com.ruoyi.ccdi.project.mapper.CcdiProjectMapper;
import com.ruoyi.ccdi.project.service.ICcdiBankTagService;
import com.ruoyi.ccdi.project.service.ICcdiProjectService;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.common.exception.ServiceException;
import com.ruoyi.lsfx.client.LsfxAnalysisClient;
import com.ruoyi.lsfx.constants.LsfxConstants;
@@ -101,6 +102,9 @@ class CcdiFileUploadServiceImplTest {
@Mock
private ICcdiProjectService projectService;
@Mock
private IProjectCounterpartyEnterpriseService counterpartyEnterpriseService;
@TempDir
Path tempDir;
@@ -695,7 +699,7 @@ class CcdiFileUploadServiceImplTest {
when(bankStatementMapper.countMatchedStaffCountByProjectId(PROJECT_ID)).thenReturn(2);
when(recordMapper.updateById(any(CcdiFileUploadRecord.class))).thenReturn(1);
String result = service.deleteFileUploadRecord(RECORD_ID, 9527L);
String result = service.deleteFileUploadRecord(RECORD_ID, CALLER);
assertEquals("删除成功已开始项目重新打标", result);
verify(lsfxClient, never()).deleteFiles(any(), any());
@@ -704,6 +708,7 @@ class CcdiFileUploadServiceImplTest {
RECORD_ID.equals(item.getId()) && "deleted".equals(item.getFileStatus())
));
verify(bankTagService).submitAutoRebuild(PROJECT_ID, TriggerType.AUTO_FILE_DELETE);
verify(counterpartyEnterpriseService).submitReconcile(CALLER, PROJECT_ID);
verify(projectMapper).updateById(org.mockito.ArgumentMatchers.<CcdiProject>argThat(item ->
PROJECT_ID.equals(item.getProjectId()) && Integer.valueOf(2).equals(item.getTargetCount())
));
@@ -716,7 +721,7 @@ class CcdiFileUploadServiceImplTest {
when(recordMapper.selectById(RECORD_ID)).thenReturn(record);
RuntimeException exception = assertThrows(RuntimeException.class,
() -> service.deleteFileUploadRecord(RECORD_ID, 9527L));
() -> service.deleteFileUploadRecord(RECORD_ID, CALLER));
assertTrue(exception.getMessage().contains("仅支持删除解析成功文件"));
}
@@ -729,7 +734,7 @@ class CcdiFileUploadServiceImplTest {
when(recordMapper.selectById(RECORD_ID)).thenReturn(record);
ServiceException exception = assertThrows(ServiceException.class,
() -> service.deleteFileUploadRecord(RECORD_ID, 9527L));
() -> service.deleteFileUploadRecord(RECORD_ID, CALLER));
assertTrue(exception.getMessage().contains("历史导入文件不支持删除"));
verify(lsfxClient, never()).deleteFiles(any(), any());
@@ -744,7 +749,7 @@ class CcdiFileUploadServiceImplTest {
when(recordMapper.selectById(RECORD_ID)).thenReturn(record);
when(recordMapper.updateById(any(CcdiFileUploadRecord.class))).thenReturn(1);
String result = service.deleteFileUploadRecord(RECORD_ID, 9527L);
String result = service.deleteFileUploadRecord(RECORD_ID, CALLER);
assertEquals("删除成功已开始项目重新打标", result);
verify(lsfxClient, never()).deleteFiles(any(), any());

View File

@@ -9,6 +9,8 @@ import com.ruoyi.ccdi.project.mapper.CcdiBankStatementMapper;
import com.ruoyi.ccdi.project.mapper.CcdiFileUploadRecordMapper;
import com.ruoyi.ccdi.project.mapper.CcdiProjectMapper;
import com.ruoyi.ccdi.project.service.ICcdiBankTagService;
import com.ruoyi.ccdi.project.service.IProjectCounterpartyEnterpriseService;
import com.ruoyi.lsfx.domain.CallerContext;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
@@ -34,6 +36,8 @@ import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
class CcdiProjectHistoryImportServiceImplTest {
private static final CallerContext CALLER = CallerContext.of(9527L, "tester");
@InjectMocks
private CcdiProjectHistoryImportServiceImpl service;
@@ -52,6 +56,9 @@ class CcdiProjectHistoryImportServiceImplTest {
@Mock
private ICcdiBankTagService bankTagService;
@Mock
private IProjectCounterpartyEnterpriseService counterpartyEnterpriseService;
@Test
void shouldFilterStatementsByTrxDateAndDeduplicateAcrossSourceProjects() {
CcdiProjectImportHistoryDTO dto = buildImportDto();
@@ -80,7 +87,7 @@ class CcdiProjectHistoryImportServiceImplTest {
when(projectMapper.selectById(11L)).thenReturn(buildProject(11L, "历史项目A"));
when(projectMapper.selectById(12L)).thenReturn(buildProject(12L, "历史项目B"));
service.submitImport(90L, 3001, dto, "tester");
service.submitImport(90L, 3001, dto, CALLER);
assertEquals(2, insertedStatements.get().size());
assertTrue(insertedStatements.get().stream().allMatch(item -> Long.valueOf(90L).equals(item.getProjectId())));
@@ -120,7 +127,7 @@ class CcdiProjectHistoryImportServiceImplTest {
when(projectMapper.selectById(11L)).thenReturn(buildProject(11L, "历史项目A"));
service.submitImport(90L, 3001, dto, "tester");
service.submitImport(90L, 3001, dto, CALLER);
assertEquals(1, insertedStatements.get().size());
assertNotEquals(101, insertedStatements.get().get(0).getBatchId());
@@ -148,7 +155,7 @@ class CcdiProjectHistoryImportServiceImplTest {
when(projectMapper.selectById(90L)).thenReturn(buildProject(90L, "新项目"));
when(bankStatementMapper.countMatchedStaffCountByProjectId(90L)).thenReturn(3);
service.submitImport(90L, 3001, dto, "tester");
service.submitImport(90L, 3001, dto, CALLER);
verify(projectMapper).updateById(org.mockito.ArgumentMatchers.<CcdiProject>argThat(project ->
Long.valueOf(90L).equals(project.getProjectId()) && Integer.valueOf(3).equals(project.getTargetCount())

View File

@@ -309,7 +309,7 @@ class CcdiProjectServiceImplTest {
(CcdiProjectHistoryImportSubmittedEvent) eventCaptor.getValue();
assertEquals(90L, event.getTargetProjectId());
assertEquals(3001, event.getTargetLsfxProjectId());
assertEquals("tester", event.getOperator());
assertEquals(CALLER, event.getCaller());
assertEquals(dto, event.getDto());
} finally {
TransactionSynchronizationManager.clearSynchronization();