Merge pull request 'kafka producer和consumer初始操作修改' (#101) from DavidZeng/gitlink-notification-system:dev_add_middleware_conf into master

This commit is contained in:
baladiwei 2021-09-18 14:18:27 +08:00
commit 94a7686c4d
6 changed files with 47 additions and 42 deletions

View File

@ -33,12 +33,17 @@ public class KafkaProducerConfig {
@Bean
ProducerFactory<String, String> producerConfigs() {
Map<String, Object> props = new HashMap<>();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.RETRIES_CONFIG, retries);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
// 如果当前配置信息内没有kafka producer相关配置则不做参数设置此处解决没有设置配置信息而引起的空指针异常
if (null != bootstrapServers) {
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
props.put(ProducerConfig.RETRIES_CONFIG, retries);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, batchSize);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
}
return new DefaultKafkaProducerFactory<>(props);
}
}

View File

@ -14,10 +14,7 @@ import org.springframework.util.concurrent.ListenableFuture;
import org.springframework.util.concurrent.ListenableFutureCallback;
import javax.annotation.PostConstruct;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.*;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
@ -37,6 +34,18 @@ public class KafkaUtil {
@Value("${spring.kafka.producer.bootstrap_servers:#{null}}")
private String kafkaServer;
@Value("${spring.kafka.producer.topic:#{null}")
private String topic;
@Value("${spring.kafka.producer.mail_topic:#{null}")
private String mailTopic;
@Value("${spring.kafka.producer.partitions:#{null}}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor:#{null}}")
private Short replicationFactor;
private AdminClient adminClient;
@Autowired
@ -46,16 +55,38 @@ public class KafkaUtil {
* 初始化AdminClient
* '@PostConstruct该注解被用来修饰一个非静态的void方法
* @PostConstruct修饰的方法会在服务器加载Servlet的时候运行并且只会被服务器执行一次
* PostConstruct在构造函数之后执行init方法之前执行
* PostConstruct在构造函数之后执行init方法之前执行ls
*/
@PostConstruct
private void initAdminClient() {
Map<String, Object> props = new HashMap<>(1);
// 如果当前配置信息内没有kafka producer相关配置则不对adminClient做初始化
if (kafkaServer == null) {
return;
}
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServer);
adminClient = KafkaAdminClient.create(props);
// 初始化topics
// TODO: 如果后面新增topic需要在这里添加
if (null !=partitions && null != replicationFactor) {
List<NewTopic> topics = new LinkedList<>();
if (null != topic) {
topics.add(new NewTopic(topic, partitions, replicationFactor));
}
if (null != mailTopic) {
topics.add(new NewTopic(mailTopic, partitions, replicationFactor));
}
if (!topics.isEmpty()) {
this.createTopic(topics);
}
}
}
public void createTopic(Collection<NewTopic> newTopics) {

View File

@ -29,12 +29,6 @@ public class EmailJobsListener {
@Value("${spring.kafka.producer.topic_new_email_remind}")
private String gitlinkNewEmailRemindTopic;
@Value("${spring.kafka.producer.partitions}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor}")
private Short replicationFactor;
@Autowired
private KafkaUtil kafkaUtil;
@ -45,7 +39,6 @@ public class EmailJobsListener {
Boolean flag = emailJobsService.sendEmail(newEmailJobVo);
//if the message is inserted successfully, send a new email-job message to kafka
if (flag){
kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkNewEmailRemindTopic, partitions, replicationFactor)));
kafkaUtil.sendMessage(gitlinkNewEmailRemindTopic, JSONObject.toJSONString(newEmailJobVo));
}
} catch (Exception e) {

View File

@ -18,12 +18,6 @@ spring:
min-idle: 0
timeout: 1000
kafka:
producer:
bootstrap_servers: kafka1:9092,kafka2:9092
retries: 5
batch_size: 16384
datasource:
driver-class-name: com.mysql.jdbc.Driver
url: jdbc:mysql://mysql:3306/gitlink_notification?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&serverTimezone=GMT%2B8&allowMultiQueries=true&useSSL=false

View File

@ -10,7 +10,6 @@ import cn.org.gitlink.notification.model.dao.entity.vo.NewEmailJobVo;
import com.alibaba.fastjson.JSONObject;
import io.swagger.annotations.ApiOperation;
import io.swagger.annotations.ApiParam;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.BeanUtils;
@ -22,7 +21,6 @@ import org.springframework.validation.BindingResult;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
import java.util.Arrays;
import java.util.Map;
@RestController
@ -35,12 +33,6 @@ public class EmailJobsController {
@Value("${spring.kafka.producer.topic_email}")
private String gitlinkEmailTopic;
@Value("${spring.kafka.producer.partitions}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor}")
private Short replicationFactor;
@Autowired
private KafkaUtil kafkaUtil;
@ -72,7 +64,6 @@ public class EmailJobsController {
newEmailJobVo.setPlatform(platform);
try {
kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkEmailTopic, partitions, replicationFactor)));
kafkaUtil.sendMessage(gitlinkEmailTopic, JSONObject.toJSONString(newEmailJobVo));
return DataPacketUtil.jsonSuccessResult();
} catch (Exception e) {

View File

@ -13,7 +13,6 @@ import cn.org.gitlink.notification.model.service.notification.SysNotificationSer
import com.alibaba.fastjson.JSONObject;
import io.swagger.annotations.ApiOperation;
import io.swagger.annotations.ApiParam;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.springframework.beans.BeanUtils;
@ -25,7 +24,6 @@ import org.springframework.validation.BindingResult;
import org.springframework.validation.annotation.Validated;
import org.springframework.web.bind.annotation.*;
import java.util.Arrays;
import java.util.Map;
@RestController
@ -36,12 +34,6 @@ public class NotificationController {
@Value("${spring.kafka.producer.topic}")
private String gitlinkNotificationTopic;
@Value("${spring.kafka.producer.partitions}")
private Integer partitions;
@Value("${spring.kafka.producer.replication_factor}")
private Short replicationFactor;
@Autowired
private KafkaUtil kafkaUtil;
@ -79,7 +71,6 @@ public class NotificationController {
newSysNotificationVo.setPlatform(platform);
try {
kafkaUtil.createTopic(Arrays.asList(new NewTopic(gitlinkNotificationTopic, partitions, replicationFactor)));
kafkaUtil.sendMessage(gitlinkNotificationTopic, JSONObject.toJSONString(newSysNotificationVo));
return DataPacketUtil.jsonSuccessResult();
} catch (Exception e) {