commit 047ae94f6490d661357178ca87da1f13c0e35394
Author: zhouxiunai <154707516@qq.com>
Date: Fri Nov 23 09:08:15 2018 +0800
首次提交代码
diff --git a/.gitignore b/.gitignore
new file mode 100644
index 00000000..a00f7f64
--- /dev/null
+++ b/.gitignore
@@ -0,0 +1,91 @@
+logs/*
+
+# SBT
+dist/*
+boot/
+project/boot/
+project/plugins/project/
+lib_managed/
+src_managed/
+target/
+.history
+.mvn/
+mvnw
+mvnw.cmd
+
+# IntelliJ
+.idea/
+*.iml
+*.ipr
+*.iws
+out/
+
+# Eclipse
+.cache
+.classpath
+.loadpath
+.metadata
+.project
+.scala_dependencies
+.settings
+.target/
+*.launch
+
+# NetBeans
+nbproject/
+nbbuild/
+dist/
+nbdist/
+nbactions.xml
+nb-configuration.xml
+
+# TextMate
+*.tmproj
+*.tmproject
+tmtags
+
+# Sublime Text
+*.sublime-workspace
+
+# vim
+.*.s[a-w][a-z]
+*.un~
+Session.vim
+.netrwhist
+*~
+
+# Emacs
+*~
+\#*\#
+/.emacs.desktop
+/.emacs.desktop.lock
+.elc
+auto-save-list
+tramp
+.\#*
+.org-id-locations
+*_archive
+
+# Mac OS X
+.DS_Store
+.AppleDouble
+.LSOverride
+Icon
+._*
+.Spotlight-V100
+.Trashes
+
+# Windows
+Thumbs.db
+ehthumbs.db
+Desktop.ini
+$RECYCLE.BIN/
+
+## Gradle ###
+.gradle
+build/
+# Ignore Gradle GUI config
+gradle-app.setting
+# Avoid ignoring Gradle wrapper jar file (.jar files are usually ignored)
+!gradle-wrapper.jar
+/.factorypath
diff --git a/.mvn/wrapper/maven-wrapper.jar b/.mvn/wrapper/maven-wrapper.jar
new file mode 100644
index 00000000..01e67997
Binary files /dev/null and b/.mvn/wrapper/maven-wrapper.jar differ
diff --git a/.mvn/wrapper/maven-wrapper.properties b/.mvn/wrapper/maven-wrapper.properties
new file mode 100644
index 00000000..71793467
--- /dev/null
+++ b/.mvn/wrapper/maven-wrapper.properties
@@ -0,0 +1 @@
+distributionUrl=https://repo.maven.apache.org/maven2/org/apache/maven/apache-maven/3.5.4/apache-maven-3.5.4-bin.zip
diff --git a/pom.xml b/pom.xml
new file mode 100644
index 00000000..ab56ab75
--- /dev/null
+++ b/pom.xml
@@ -0,0 +1,78 @@
+
+
+ 4.0.0
+
+ com.gzzn.omms
+ msgexchange-api
+ 0.0.1-SNAPSHOT
+ jar
+
+ msgexchange-api
+ msg exchange service
+
+
+ org.springframework.boot
+ spring-boot-starter-parent
+ 1.5.17.RELEASE
+
+
+
+
+ UTF-8
+ UTF-8
+ 1.8
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-web
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+
+ org.springframework.kafka
+ spring-kafka
+
+
+
+ org.springframework.boot
+ spring-boot-starter-data-jpa
+
+
+
+ org.springframework.boot
+ spring-boot-starter-jdbc
+
+
+
+ com.fasterxml.jackson.dataformat
+ jackson-dataformat-xml
+
+
+
+ com.oracle
+ ojdbc6
+ 11.2.0.3
+ system
+ ${project.basedir}/src/main/resources/libs/ojdbc6-11.2.0.3.jar
+
+
+
+
+
+
+ org.springframework.boot
+ spring-boot-maven-plugin
+
+
+
+
+
+
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplication.java b/src/main/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplication.java
new file mode 100644
index 00000000..6f1525ef
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplication.java
@@ -0,0 +1,14 @@
+package com.gzzn.omms.msgexchangeapi;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.scheduling.annotation.EnableScheduling;
+
+@SpringBootApplication
+@EnableScheduling
+public class MsgexchangeApiApplication {
+
+ public static void main(String[] args) {
+ SpringApplication.run(MsgexchangeApiApplication.class, args);
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/config/KafkaProducerConfig.java b/src/main/java/com/gzzn/omms/msgexchangeapi/config/KafkaProducerConfig.java
new file mode 100644
index 00000000..2f80f7ea
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/config/KafkaProducerConfig.java
@@ -0,0 +1,52 @@
+package com.gzzn.omms.msgexchangeapi.config;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.kafka.annotation.EnableKafka;
+import org.springframework.kafka.core.DefaultKafkaProducerFactory;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.kafka.core.ProducerFactory;
+
+@Configuration
+@EnableKafka
+public class KafkaProducerConfig {
+
+ @Value("${kafka.producer.servers}")
+ private String servers;
+ @Value("${kafka.producer.retries}")
+ private int retries;
+ @Value("${kafka.producer.batch.size}")
+ private int batchSize;
+ @Value("${kafka.producer.linger}")
+ private int linger;
+ @Value("${kafka.producer.buffer.memory}")
+ private int bufferMemory;
+
+
+ public Map producerConfigs() {
+ Map props = new HashMap<>();
+ props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
+ props.put(ProducerConfig.RETRIES_CONFIG, retries);
+ props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
+ props.put(ProducerConfig.LINGER_MS_CONFIG, linger);
+ props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, bufferMemory);
+ props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
+ props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
+ return props;
+ }
+
+ public ProducerFactory producerFactory() {
+ return new DefaultKafkaProducerFactory<>(producerConfigs());
+ }
+
+ @Bean
+ public KafkaTemplate kafkaTemplate() {
+ return new KafkaTemplate(producerFactory());
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/controller/KafkaController.java b/src/main/java/com/gzzn/omms/msgexchangeapi/controller/KafkaController.java
new file mode 100644
index 00000000..ba3168fb
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/controller/KafkaController.java
@@ -0,0 +1,22 @@
+package com.gzzn.omms.msgexchangeapi.controller;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.web.bind.annotation.GetMapping;
+import org.springframework.web.bind.annotation.RequestMapping;
+import org.springframework.web.bind.annotation.RestController;
+
+import com.gzzn.omms.msgexchangeapi.dto.ResponseDto;
+import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
+
+@RestController
+@RequestMapping("/kafka")
+public class KafkaController {
+ @Autowired
+ private IKafkaService kafkaservice;
+
+ @GetMapping(value = "/send")
+ public ResponseDto sendKafka(String msg,String topic) {
+ kafkaservice.msgSend(topic, msg);
+ return ResponseDto.success();
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java b/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java
new file mode 100644
index 00000000..3c73576c
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/dao/CminmsgDao.java
@@ -0,0 +1,34 @@
+package com.gzzn.omms.msgexchangeapi.dao;
+
+import java.util.List;
+
+import org.springframework.data.domain.Pageable;
+import org.springframework.data.jpa.repository.Query;
+import org.springframework.data.repository.CrudRepository;
+
+import com.gzzn.omms.msgexchangeapi.entiy.Cminmsg;
+
+
+public interface CminmsgDao extends CrudRepository {
+
+ /**
+ *
+ * @return
+ */
+ @Query("select count(*) from Cminmsg c where cminmsgsDateProcessed is not null")
+ public Long getCminmsgsDateProcessedIsNotNullCount();
+
+
+ /**
+ * 根据处理时间字段查询列表,Null情况
+ * @param date
+ * @return
+ */
+ public List findByCminmsgsDateProcessedIsNull(Pageable pageable);
+
+ /**
+ * 非空情况
+ * @return
+ */
+ public List findByCminmsgsDateProcessedIsNotNull(Pageable pageable);
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/dto/ResponseDto.java b/src/main/java/com/gzzn/omms/msgexchangeapi/dto/ResponseDto.java
new file mode 100644
index 00000000..54710a71
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/dto/ResponseDto.java
@@ -0,0 +1,97 @@
+package com.gzzn.omms.msgexchangeapi.dto;
+
+import com.gzzn.omms.msgexchangeapi.enums.ResultCode;
+
+/**
+ *
+ * @author xiunai
+ *
+ * @date 2018年5月16日
+ *
+ *
+ */
+public class ResponseDto {
+
+ private Boolean is_success;
+
+ private Integer err_code;
+
+ private String err_msg;
+
+ private T body;
+
+ public ResponseDto() {}
+
+ public ResponseDto(Integer code, String msg) {
+ this.setIs_success(false);
+ this.setErr_code(code);
+ this.setErr_msg(msg);
+ }
+
+ public static ResponseDto success() {
+ ResponseDto result = new ResponseDto();
+ result.setIs_success(true);
+ result.setResultCode(ResultCode.SUCCESS);
+ return result;
+ }
+
+ public static ResponseDto success(T data) {
+ ResponseDto result = new ResponseDto();
+ result.setIs_success(true);
+ result.setResultCode(ResultCode.SUCCESS);
+ result.setBody(data);
+ return result;
+ }
+
+ public static ResponseDto failure(ResultCode resultCode) {
+ ResponseDto result = new ResponseDto();
+ result.setIs_success(false);
+ result.setResultCode(resultCode);
+ return result;
+ }
+
+ public static ResponseDto failure(ResultCode resultCode, T data) {
+ ResponseDto result = new ResponseDto();
+ result.setIs_success(false);
+ result.setResultCode(resultCode);
+ result.setBody(data);
+ return result;
+ }
+
+ public void setResultCode(ResultCode code) {
+ this.setErr_code(code.code());
+ this.setErr_msg(code.message());
+ }
+
+ public T getBody() {
+ return body;
+ }
+
+ public void setBody(T body) {
+ this.body = body;
+ }
+
+ public Boolean getIs_success() {
+ return is_success;
+ }
+
+ public void setIs_success(Boolean is_success) {
+ this.is_success = is_success;
+ }
+
+ public Integer getErr_code() {
+ return err_code;
+ }
+
+ public void setErr_code(Integer err_code) {
+ this.err_code = err_code;
+ }
+
+ public String getErr_msg() {
+ return err_msg;
+ }
+
+ public void setErr_msg(String err_msg) {
+ this.err_msg = err_msg;
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/entiy/Cminmsg.java b/src/main/java/com/gzzn/omms/msgexchangeapi/entiy/Cminmsg.java
new file mode 100644
index 00000000..ab51994a
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/entiy/Cminmsg.java
@@ -0,0 +1,136 @@
+package com.gzzn.omms.msgexchangeapi.entiy;
+
+import java.io.Serializable;
+import javax.persistence.*;
+import java.util.Date;
+
+
+/**
+ * The persistent class for the CMINMSGS database table.
+ *
+ */
+@Entity
+@Table(name="CMINMSGS")
+@NamedQuery(name="Cminmsg.findAll", query="SELECT c FROM Cminmsg c")
+public class Cminmsg implements Serializable {
+ private static final long serialVersionUID = 1L;
+
+ @Id
+ @Column(name="CMINMSGS_ID")
+ private long cminmsgsId;
+
+ @Lob
+ @Column(name="CMINMSGS_CLOB_MSG")
+ private String cminmsgsClobMsg;
+
+ @Temporal(TemporalType.DATE)
+ @Column(name="CMINMSGS_DATE_PROCESSED")
+ private Date cminmsgsDateProcessed;
+
+ @Temporal(TemporalType.DATE)
+ @Column(name="CMINMSGS_DATE_RECEIVED")
+ private Date cminmsgsDateReceived;
+
+ @Column(name="CMINMSGS_STATUS")
+ private String cminmsgsStatus;
+
+ @Temporal(TemporalType.DATE)
+ @Column(name="CMINMSGS_SUBSYSTEM_DATE_SENT")
+ private Date cminmsgsSubsystemDateSent;
+
+ @Column(name="CMINMSGS_SUBSYSTEM_NAME")
+ private String cminmsgsSubsystemName;
+
+ @Column(name="CMINMSGS_SUBSYSTEM_SEQUENCE")
+ private String cminmsgsSubsystemSequence;
+
+ @Column(name="CMINMSGS_SUBSYSTEM_SUBTYPE")
+ private String cminmsgsSubsystemSubtype;
+
+ @Column(name="CMINMSGS_SUBSYSTEM_TYPE")
+ private String cminmsgsSubsystemType;
+
+ public Cminmsg() {
+ }
+
+ public long getCminmsgsId() {
+ return this.cminmsgsId;
+ }
+
+ public void setCminmsgsId(long cminmsgsId) {
+ this.cminmsgsId = cminmsgsId;
+ }
+
+ public String getCminmsgsClobMsg() {
+ return this.cminmsgsClobMsg;
+ }
+
+ public void setCminmsgsClobMsg(String cminmsgsClobMsg) {
+ this.cminmsgsClobMsg = cminmsgsClobMsg;
+ }
+
+ public Date getCminmsgsDateProcessed() {
+ return this.cminmsgsDateProcessed;
+ }
+
+ public void setCminmsgsDateProcessed(Date cminmsgsDateProcessed) {
+ this.cminmsgsDateProcessed = cminmsgsDateProcessed;
+ }
+
+ public Date getCminmsgsDateReceived() {
+ return this.cminmsgsDateReceived;
+ }
+
+ public void setCminmsgsDateReceived(Date cminmsgsDateReceived) {
+ this.cminmsgsDateReceived = cminmsgsDateReceived;
+ }
+
+ public String getCminmsgsStatus() {
+ return this.cminmsgsStatus;
+ }
+
+ public void setCminmsgsStatus(String cminmsgsStatus) {
+ this.cminmsgsStatus = cminmsgsStatus;
+ }
+
+ public Date getCminmsgsSubsystemDateSent() {
+ return this.cminmsgsSubsystemDateSent;
+ }
+
+ public void setCminmsgsSubsystemDateSent(Date cminmsgsSubsystemDateSent) {
+ this.cminmsgsSubsystemDateSent = cminmsgsSubsystemDateSent;
+ }
+
+ public String getCminmsgsSubsystemName() {
+ return this.cminmsgsSubsystemName;
+ }
+
+ public void setCminmsgsSubsystemName(String cminmsgsSubsystemName) {
+ this.cminmsgsSubsystemName = cminmsgsSubsystemName;
+ }
+
+ public String getCminmsgsSubsystemSequence() {
+ return this.cminmsgsSubsystemSequence;
+ }
+
+ public void setCminmsgsSubsystemSequence(String cminmsgsSubsystemSequence) {
+ this.cminmsgsSubsystemSequence = cminmsgsSubsystemSequence;
+ }
+
+ public String getCminmsgsSubsystemSubtype() {
+ return this.cminmsgsSubsystemSubtype;
+ }
+
+ public void setCminmsgsSubsystemSubtype(String cminmsgsSubsystemSubtype) {
+ this.cminmsgsSubsystemSubtype = cminmsgsSubsystemSubtype;
+ }
+
+ public String getCminmsgsSubsystemType() {
+ return this.cminmsgsSubsystemType;
+ }
+
+ public void setCminmsgsSubsystemType(String cminmsgsSubsystemType) {
+ this.cminmsgsSubsystemType = cminmsgsSubsystemType;
+ }
+
+}
\ No newline at end of file
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/enums/ResultCode.java b/src/main/java/com/gzzn/omms/msgexchangeapi/enums/ResultCode.java
new file mode 100644
index 00000000..2c4ab4d6
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/enums/ResultCode.java
@@ -0,0 +1,83 @@
+package com.gzzn.omms.msgexchangeapi.enums;
+
+public enum ResultCode {
+ /* 成功状态码 */
+ SUCCESS(1, "成功"),
+ Failure(0,"失败"),
+
+ /* 参数错误:10001-19999 */
+ PARAM_IS_INVALID(10001, "参数无效"),
+ PARAM_IS_BLANK(10002, "参数为空"),
+ PARAM_TYPE_BIND_ERROR(10003, "参数类型错误"),
+ PARAM_NOT_COMPLETE(10004, "参数缺失"),
+
+ /* 用户错误:20001-29999*/
+ USER_NOT_LOGGED_IN(20001, "用户未登录"),
+ USER_LOGIN_ERROR(20002, "账号不存在或密码错误"),
+ USER_ACCOUNT_FORBIDDEN(20003, "账号已被禁用"),
+ USER_NOT_EXIST(20004, "用户不存在"),
+ USER_HAS_EXISTED(20005, "用户已存在"),
+
+ /* 业务错误:30001-39999 */
+ SPECIFIED_QUESTIONED_USER_NOT_EXIST(30001, "某业务出现问题"),
+
+ /* 系统错误:40001-49999 */
+ SYSTEM_INNER_ERROR(40001, "系统繁忙,请稍后重试"),
+
+ /* 数据错误:50001-599999 */
+ RESULE_DATA_NONE(50001, "数据未找到"),
+ DATA_IS_WRONG(50002, "数据有误"),
+ DATA_ALREADY_EXISTED(50003, "数据已存在"),
+
+ /* 接口错误:60001-69999 */
+ INTERFACE_INNER_INVOKE_ERROR(60001, "内部系统接口调用异常"),
+ INTERFACE_OUTTER_INVOKE_ERROR(60002, "外部系统接口调用异常"),
+ INTERFACE_FORBID_VISIT(60003, "该接口禁止访问"),
+ INTERFACE_ADDRESS_INVALID(60004, "接口地址无效"),
+ INTERFACE_REQUEST_TIMEOUT(60005, "接口请求超时"),
+ INTERFACE_EXCEED_LOAD(60006, "接口负载过高"),
+ INTERFACE_NOT_IMPLEMENT(60007, "接口暂未实现"),
+
+ /* 权限错误:70001-79999 */
+ PERMISSION_NO_ACCESS(70001, "无访问权限");
+
+
+ private Integer code;
+ private String message;
+
+ ResultCode(Integer code, String message) {
+ this.code = code;
+ this.message = message;
+ }
+
+ public Integer code() {
+ return this.code;
+ }
+
+ public String message() {
+ return this.message;
+ }
+
+ public static String getMessage(String name) {
+ for (ResultCode item : ResultCode.values()) {
+ if (item.name().equals(name)) {
+ return item.message;
+ }
+ }
+ return name;
+ }
+
+ public static Integer getCode(String name) {
+ for (ResultCode item : ResultCode.values()) {
+ if (item.name().equals(name)) {
+ return item.code;
+ }
+ }
+ return null;
+ }
+
+ @Override
+ public String toString() {
+ return this.name();
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImpl.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImpl.java
new file mode 100644
index 00000000..f8d2e714
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImpl.java
@@ -0,0 +1,28 @@
+package com.gzzn.omms.msgexchangeapi.service;
+
+import java.io.IOException;
+import java.util.Map;
+
+import org.springframework.stereotype.Service;
+
+import com.fasterxml.jackson.core.JsonParseException;
+import com.fasterxml.jackson.databind.JsonMappingException;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.dataformat.xml.XmlMapper;
+
+@Service
+public class ExchangeServiceImpl implements IExchangeService {
+
+ @Override
+ public String xmlToJson(String xml) throws JsonParseException, JsonMappingException, IOException {
+
+ ObjectMapper xmlMapper = new XmlMapper();
+ Map map = xmlMapper.readValue(xml, Map.class);
+
+ ObjectMapper jsonMapper = new ObjectMapper();
+ String json = jsonMapper.writeValueAsString(map);
+
+ return json;
+ }
+
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/IExchangeService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IExchangeService.java
new file mode 100644
index 00000000..a1b0209b
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IExchangeService.java
@@ -0,0 +1,16 @@
+package com.gzzn.omms.msgexchangeapi.service;
+
+import java.io.IOException;
+
+import com.fasterxml.jackson.core.JsonParseException;
+import com.fasterxml.jackson.databind.JsonMappingException;
+
+/**
+ * 消息转换服务
+ * @author Administrator
+ *
+ */
+public interface IExchangeService {
+
+ public String xmlToJson(String xml)throws JsonParseException, JsonMappingException, IOException;
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/IKafkaService.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IKafkaService.java
new file mode 100644
index 00000000..88c8c659
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/IKafkaService.java
@@ -0,0 +1,5 @@
+package com.gzzn.omms.msgexchangeapi.service;
+
+public interface IKafkaService {
+ public void msgSend(String topic,String msg) ;
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java b/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java
new file mode 100644
index 00000000..875c5508
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/service/KafkaServiceImpl.java
@@ -0,0 +1,18 @@
+package com.gzzn.omms.msgexchangeapi.service;
+
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.kafka.core.KafkaTemplate;
+import org.springframework.stereotype.Service;
+
+@Service
+public class KafkaServiceImpl implements IKafkaService{
+
+ @Autowired
+ private KafkaTemplate kafkaTemplate;
+
+
+ @Override
+ public void msgSend(String topic, String msg) {
+ kafkaTemplate.send(topic, msg);
+ }
+}
diff --git a/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java b/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java
new file mode 100644
index 00000000..168f2814
--- /dev/null
+++ b/src/main/java/com/gzzn/omms/msgexchangeapi/task/ExchangeTask.java
@@ -0,0 +1,76 @@
+package com.gzzn.omms.msgexchangeapi.task;
+
+import java.io.IOException;
+import java.util.List;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.data.domain.PageRequest;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+import com.fasterxml.jackson.core.JsonParseException;
+import com.fasterxml.jackson.databind.JsonMappingException;
+import com.gzzn.omms.msgexchangeapi.dao.CminmsgDao;
+import com.gzzn.omms.msgexchangeapi.entiy.Cminmsg;
+import com.gzzn.omms.msgexchangeapi.service.IExchangeService;
+import com.gzzn.omms.msgexchangeapi.service.IKafkaService;
+
+/**
+ * 转换任务,获取cim的消息,转换成json,写入到kafka队列中
+ * @author Administrator
+ *
+ */
+@Component
+public class ExchangeTask {
+ private static Logger logger = LoggerFactory.getLogger(ExchangeTask.class);
+
+ @Autowired
+ private IKafkaService kafkaservice;
+
+ @Autowired
+ IExchangeService exchangeService;
+
+ @Autowired
+ private CminmsgDao cminmsgDao;
+
+ @Scheduled(cron="0 0/1 * * * ?")
+ public void corn()
+ {
+ logger.info("定时任务启动....");
+
+ Long maxPerLong = 1000L; //每次最大发送多少条消息
+ Long lenTotal = cminmsgDao.getCminmsgsDateProcessedIsNotNullCount();
+ Long lenSend = 0L;
+ Integer pageIndex = 0;
+ Integer pageSize = 50;
+ while(lenSend < lenTotal && lenSend < maxPerLong)
+ {
+ List ls = cminmsgDao.findByCminmsgsDateProcessedIsNotNull(new PageRequest(pageIndex,pageSize));
+ ls.forEach(x -> {
+ String strJson = null;
+ try {
+ strJson = exchangeService.xmlToJson(x.getCminmsgsClobMsg());
+ if(!StringUtils.isEmpty(strJson))
+ {
+ kafkaservice.msgSend("topic2", strJson);
+ }
+ } catch (JsonParseException e) {
+ e.printStackTrace();
+ } catch (JsonMappingException e) {
+ e.printStackTrace();
+ } catch (IOException e) {
+ e.printStackTrace();
+ }
+ });
+
+
+ lenSend += ls.size();
+ pageIndex++;
+ }
+ //
+ logger.info("同步完成");
+ }
+}
diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml
new file mode 100644
index 00000000..33f78422
--- /dev/null
+++ b/src/main/resources/application-dev.yml
@@ -0,0 +1,40 @@
+kafka:
+ producer:
+ retries: 0
+ servers: 130.120.3.237:32796,130.120.3.237:32797,130.120.3.237:32798
+ linger: 1
+ batch:
+ size: 4096
+ buffer:
+ memory: 40960
+ max:
+ request:
+ size: 10240000
+ consumer:
+ auto:
+ offset:
+ reset: latest
+ commit:
+ interval: 100
+ servers: 130.120.3.237:32796,130.120.3.237:32797,130.120.3.237:32798
+ zookeeper:
+ connect: 130.120.3.237:2181
+ session:
+ timeout: 6000
+ enable:
+ auto:
+ commit: true
+ topic: test
+ concurrency: 10
+ group:
+ id: test
+spring:
+ datasource:
+ username: ommsxc
+ password: ommsxc
+ url: jdbc:oracle:thin:@130.120.3.236:1521/xe
+ driver: oracle.jdbc.driver.OracleDriver
+ tomcat:
+ max-active: 30
+ test-on-borrow: true
+ initial-size: 3
\ No newline at end of file
diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml
new file mode 100644
index 00000000..6fbec020
--- /dev/null
+++ b/src/main/resources/application.yml
@@ -0,0 +1,3 @@
+spring:
+ profiles:
+ active: dev
\ No newline at end of file
diff --git a/src/main/resources/libs/ojdbc6-11.2.0.3.jar b/src/main/resources/libs/ojdbc6-11.2.0.3.jar
new file mode 100644
index 00000000..01da074d
Binary files /dev/null and b/src/main/resources/libs/ojdbc6-11.2.0.3.jar differ
diff --git a/src/test/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplicationTests.java b/src/test/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplicationTests.java
new file mode 100644
index 00000000..cb7a8b95
--- /dev/null
+++ b/src/test/java/com/gzzn/omms/msgexchangeapi/MsgexchangeApiApplicationTests.java
@@ -0,0 +1,16 @@
+package com.gzzn.omms.msgexchangeapi;
+
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.test.context.junit4.SpringRunner;
+
+@RunWith(SpringRunner.class)
+@SpringBootTest
+public class MsgexchangeApiApplicationTests {
+
+ @Test
+ public void contextLoads() {
+ }
+
+}
diff --git a/src/test/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImplTest.java b/src/test/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImplTest.java
new file mode 100644
index 00000000..bbf04e3c
--- /dev/null
+++ b/src/test/java/com/gzzn/omms/msgexchangeapi/service/ExchangeServiceImplTest.java
@@ -0,0 +1,47 @@
+package com.gzzn.omms.msgexchangeapi.service;
+
+import java.io.IOException;
+
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.test.context.junit4.SpringRunner;
+
+import com.fasterxml.jackson.core.JsonParseException;
+import com.fasterxml.jackson.databind.JsonMappingException;
+
+@RunWith(SpringRunner.class)
+@SpringBootTest
+public class ExchangeServiceImplTest {
+ @Autowired
+ private IExchangeService exchangeService;
+
+
+ @Test
+ public void testXmlToJson() throws JsonParseException, JsonMappingException, IOException {
+ String xmlString = "\r\n" +
+ "\r\n" +
+ " \r\n" +
+ " AODB\r\n" +
+ "754508\r\n" +
+ "20181112221954\r\n" +
+ "FLOP\r\n" +
+ "PSDT\r\n" +
+ "\r\n" +
+ "\r\n" +
+ "11450149\r\n" +
+ "MU-MU5851-A-12NOV182250-D\r\n" +
+ "\r\n" +
+ "343\r\n" +
+ "12NOV182230\r\n" +
+ "12NOV182300\r\n" +
+ "\r\n" +
+ "\r\n" +
+ "";
+ String jsonString = exchangeService.xmlToJson(xmlString);
+ System.out.println(jsonString);
+
+ //assertThat(jsonPath("$.", matcher)).
+ }
+}