kafka producer和consumer初始操作修改
1.初始化时创建topic 2.初始化时没有kafka相关配置信息则不进行kafka连接
This commit is contained in:
parent
728efe9a57
commit
ddb7098080
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
Loading…
Reference in New Issue