tao.mo

kafka-cndc-server | 修改 | kafka501消息处理

mt
2024年11月1日13:46:12
......@@ -11,7 +11,7 @@ import org.springframework.web.WebApplicationInitializer;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
@SpringBootApplication(exclude = SecurityAutoConfiguration.class)
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.customer.repository.dao.*"})
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.customer.repository.dao"})
@EnableTransactionManagement
@EnableSwagger2
public class CustomerApplication extends SpringBootServletInitializer implements WebApplicationInitializer {
......
......@@ -11,7 +11,7 @@ import org.springframework.web.WebApplicationInitializer;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
@SpringBootApplication(exclude = SecurityAutoConfiguration.class)
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.kafka.repository.dao.*"})
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.kafka.repository.dao"})
@EnableTransactionManagement
@EnableSwagger2
public class KafkaApplication extends SpringBootServletInitializer implements WebApplicationInitializer {
......
package com.fedex.connect.kafka.repository.base;
import com.fedex.connect.common.dao.biz.CeInfoMapper;
import com.fedex.connect.common.dao.biz.ConsignmentMapper;
import com.fedex.connect.common.dao.sys.KafkaStorageHistoryMapper;
import com.fedex.connect.common.dao.sys.KafkaTemporaryStorageMapper;
import com.fedex.connect.kafka.repository.dao.CeInfoMapperExt;
......@@ -15,6 +16,8 @@ public class AbstractDaoRepository {
@Autowired
protected CeInfoMapper ceInfoMapper;
@Autowired
protected ConsignmentMapper consignmentMapper;
@Autowired
protected CeInfoMapperExt ceInfoMapperExt;
@Autowired
protected KafkaTemporaryStorageMapperExt kafkaTemporaryStorageMapperExt;
......
package com.fedex.connect.kafka.repository.repo;
import com.fedex.connect.common.model.biz.Consignment;
import java.util.List;
public interface IConsignmentRepository {
Consignment save(Consignment entity);
void saveAll(List<Consignment> list);
}
package com.fedex.connect.kafka.repository.repo.impl;
import com.fedex.connect.common.model.biz.Consignment;
import com.fedex.connect.kafka.repository.base.AbstractDaoRepository;
import com.fedex.connect.kafka.repository.repo.IConsignmentRepository;
import org.springframework.stereotype.Repository;
import java.util.List;
@Repository
public class ConsignmentRepositoryImpl extends AbstractDaoRepository implements IConsignmentRepository {
@Override
public Consignment save(Consignment entity) {
if(entity != null) {
if (entity.getId() == null || entity.getId().longValue() <= 0) {
consignmentMapper.insertSelective(entity);
} else if (entity.getId() != null && entity.getId().longValue() > 0){
consignmentMapper.updateByPrimaryKeySelective(entity);
}
}
return entity;
}
@Override
public void saveAll(List<Consignment> list) {
list.stream().forEach(p->{
if(p.getId() == null || p.getId().longValue() <=0){
consignmentMapper.insertSelective(p);
}else if(p.getId() != null && p.getId().longValue() >0){
consignmentMapper.updateByPrimaryKeySelective(p);
}
});
}
}
package com.fedex.connect.kafka.service.base;
import com.fedex.connect.kafka.repository.repo.ICeInfoRepository;
import com.fedex.connect.kafka.repository.repo.IConsignmentRepository;
import com.fedex.connect.kafka.repository.repo.IKafkaStorageHistoryRepository;
import com.fedex.connect.kafka.repository.repo.IKafkaTemporaryStorageRepository;
import org.springframework.beans.factory.annotation.Autowired;
......@@ -17,4 +18,6 @@ public class BaseService {
protected IKafkaTemporaryStorageRepository kafkaTemporaryStorageRepository;
@Autowired
protected ICeInfoRepository ceInfoRepository;
@Autowired
protected IConsignmentRepository consignmentRepository;
}
\ No newline at end of file
......
package com.fedex.connect.kafka.service.impl;
import com.alibaba.fastjson.JSONObject;
import com.fedex.connect.common.dependencies.util.DateUtil;
import com.fedex.connect.common.dependencies.util.Utils;
import com.fedex.connect.common.model.biz.CeInfo;
import com.fedex.connect.common.model.biz.Consignment;
import com.fedex.connect.common.model.sys.KafkaTemporaryStorage;
import com.fedex.connect.kafka.constants.Constant;
import com.fedex.connect.kafka.data.dto.DmProcessResults;
import com.fedex.connect.kafka.data.dto.dm.dm501.Dm501ConAddresses;
import com.fedex.connect.kafka.data.dto.dm.dm501.Dm501Consignments;
import com.fedex.connect.kafka.data.dto.dm.dm501.Dm501Data;
import com.fedex.connect.kafka.service.IDm501Service;
import com.fedex.connect.kafka.service.base.BaseService;
import org.apache.commons.lang3.StringUtils;
import com.fedex.connect.kafka.util.Dm501Util;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.List;
import java.util.Objects;
......@@ -27,6 +25,9 @@ import java.util.Objects;
@Service
public class Dm501ServiceImpl extends BaseService implements IDm501Service {
@Autowired
Dm501Util dm501Util;
/**
* @Author mt
* @Description 处理dm501消息
......@@ -48,8 +49,7 @@ public class Dm501ServiceImpl extends BaseService implements IDm501Service {
processResults.setConsignmentCode(consignment.getTrackingNumber());
//该运单号不存在创建时间为1个月内CE数据
if(Objects.isNull(bizCeInfo)){
CeInfo ceInfo = this.generateCeInfo(kafKaTemporaryStorage.getSendTime(),consignment);
ceInfoRepository.save(ceInfo);
CeInfo ceInfo = dm501Util.generateCeInfo(kafKaTemporaryStorage.getSendTime(),consignment);
}
});
return processResults;
......@@ -57,69 +57,13 @@ public class Dm501ServiceImpl extends BaseService implements IDm501Service {
/**
* @Author mt
* @Description 生成CE数据实体
* @Date 2024/4/22
* @param dm501Consignment
* @return com.fedex.connect.common.model.biz.CeInfo
* @Description 保存ce信息、以及运单表信息
* @Date 2024/11/1
* @param ceInfo
* @param consignment
* @return void
*/
private CeInfo generateCeInfo(Date sendTime, Dm501Consignments dm501Consignment){
CeInfo bizCeInfo = new CeInfo();
bizCeInfo.setConsignmentCode(dm501Consignment.getTrackingNumber());
bizCeInfo.setCreateTime(new Date());
bizCeInfo.setDestIataCode(dm501Consignment.getDestIataCode());
//地址信息
List<Dm501ConAddresses> addressList = dm501Consignment.getAddresses();
//没有地址信息,不需要处理,不进CE业务表
if(addressList == null){
return null;
}
//发件人
Dm501ConAddresses shipperAddress = this.queryDM501Addresses(addressList, Constant.CE_INFO_ADDRESS_KEYS.SHIPPER);
//收件人
Dm501ConAddresses recipientAddress = this.queryDM501Addresses(addressList, Constant.CE_INFO_ADDRESS_KEYS.RECIPIENT);
//没有收发件人地址信息,不需要处理,不进CE业务表
if(shipperAddress == null || recipientAddress == null){
return null;
}
//发件人国家
String shipperCountry = shipperAddress.getCountry();
//始发国(二字码)
bizCeInfo.setShipperCountry(shipperCountry);
//发件人国家不为TW,不需要处理,不进CE业务表
if(StringUtils.isEmpty(shipperCountry) || !shipperCountry.trim().equals(Constant.CE_INFO_KEYS.TW)){
return null;
}
//收件人国家
String recipientCountry = recipientAddress.getCountry();
//目的国家二字码
bizCeInfo.setRecipientCountry(recipientCountry);
//收件人国家不为CN,不需要处理,不进CE业务表
if(StringUtils.isEmpty(recipientCountry) || !recipientCountry.trim().equals(Constant.CE_INFO_KEYS.CN)){
return null;
}
bizCeInfo.setSendTime(sendTime);
String shipDate = dm501Consignment.getShipDate();
if(StringUtils.isNotEmpty(shipDate)){
bizCeInfo.setShipDate(DateUtil.parse(shipDate));
}
return bizCeInfo;
}
/**
* 根据类型获取地址信息
* SHIPPER/发件人---RECIPIENT/收件人---BROKEN/代理人
* @param type
* @return
*/
private Dm501ConAddresses queryDM501Addresses(List<Dm501ConAddresses> addressList,String type){
Dm501ConAddresses address = null;
for(int i = 0 ; i < addressList.size() ; i ++){
Dm501ConAddresses ad = addressList.get(i);
if(ad.getType().equals(type)){
address = ad;
break;
}
}
return address;
public void saveCeInfoAndConsignment(CeInfo ceInfo, Consignment consignment){
ceInfoRepository.save(ceInfo);
}
}
\ No newline at end of file
......
package com.fedex.connect.kafka.util;
import com.fedex.connect.common.dependencies.util.DateUtil;
import com.fedex.connect.common.model.biz.CeInfo;
import com.fedex.connect.kafka.constants.Constant;
import com.fedex.connect.kafka.data.dto.dm.dm501.Dm501ConAddresses;
import com.fedex.connect.kafka.data.dto.dm.dm501.Dm501Consignments;
import org.apache.commons.lang3.StringUtils;
import org.springframework.stereotype.Component;
import java.util.Date;
import java.util.List;
@Component
public class Dm501Util {
/**
* @Author mt
* @Description 生成CE数据实体
* @Date 2024/4/22
* @param dm501Consignment
* @return com.fedex.connect.common.model.biz.CeInfo
*/
public CeInfo generateCeInfo(Date sendTime, Dm501Consignments dm501Consignment){
CeInfo ceInfo = new CeInfo();
ceInfo.setConsignmentCode(dm501Consignment.getTrackingNumber());
ceInfo.setCreateTime(new Date());
ceInfo.setDestIataCode(dm501Consignment.getDestIataCode());
//地址信息
List<Dm501ConAddresses> addressList = dm501Consignment.getAddresses();
//没有地址信息,不需要处理,不进CE业务表
if(addressList == null){
return null;
}
//发件人
Dm501ConAddresses shipperAddress = this.queryDM501Addresses(addressList, Constant.CE_INFO_ADDRESS_KEYS.SHIPPER);
//收件人
Dm501ConAddresses recipientAddress = this.queryDM501Addresses(addressList, Constant.CE_INFO_ADDRESS_KEYS.RECIPIENT);
//没有收发件人地址信息,不需要处理,不进CE业务表
if(shipperAddress == null || recipientAddress == null){
return null;
}
//发件人国家
String shipperCountry = shipperAddress.getCountry();
//始发国(二字码)
ceInfo.setShipperCountry(shipperCountry);
//收件人国家
String recipientCountry = recipientAddress.getCountry();
//目的国家二字码
ceInfo.setRecipientCountry(recipientCountry);
//收件人国家不为CN,不需要处理,不进CE业务表
if(StringUtils.isEmpty(recipientCountry) || !recipientCountry.trim().equals(Constant.CE_INFO_KEYS.CN)){
return null;
}
ceInfo.setSendTime(sendTime);
String shipDate = dm501Consignment.getShipDate();
if(StringUtils.isNotEmpty(shipDate)){
ceInfo.setShipDate(DateUtil.parse(shipDate));
}
return ceInfo;
}
/**
* 根据类型获取地址信息
* SHIPPER/发件人---RECIPIENT/收件人---BROKEN/代理人
* @param type
* @return
*/
private Dm501ConAddresses queryDM501Addresses(List<Dm501ConAddresses> addressList,String type){
Dm501ConAddresses address = null;
for(int i = 0 ; i < addressList.size() ; i ++){
Dm501ConAddresses ad = addressList.get(i);
if(ad.getType().equals(type)){
address = ad;
break;
}
}
return address;
}
}
......@@ -11,7 +11,7 @@ import org.springframework.web.WebApplicationInitializer;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
@SpringBootApplication(exclude = SecurityAutoConfiguration.class)
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.manager.repository.dao.*"})
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.manager.repository.dao"})
@EnableTransactionManagement
@EnableSwagger2
public class ManagerApplication extends SpringBootServletInitializer implements WebApplicationInitializer {
......
......@@ -11,7 +11,7 @@ import org.springframework.web.WebApplicationInitializer;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
@SpringBootApplication(exclude = SecurityAutoConfiguration.class)
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.task.repository.dao.*"})
@MapperScan({"com.fedex.connect.common.dao.*.**","com.fedex.connect.task.repository.dao"})
@EnableTransactionManagement
@EnableSwagger2
public class TaskApplication extends SpringBootServletInitializer implements WebApplicationInitializer {
......