遇到了需要暫停消費(fèi)的場(chǎng)景,使用pause()方法暫停消費(fèi),resume()方法恢復(fù)消費(fèi),基于springboot的demo如下:
注意杈帐,如果需要暫停消費(fèi)的話,需要consumer 訂閱 topic 的方式必須是 Assign专钉。
assign 和 subscribe 的區(qū)別 :assign方法由用戶直接手動(dòng)consumer實(shí)例消費(fèi)哪些具體分區(qū)挑童,assign的consumer不會(huì)擁有kafka的group management機(jī)制,也就是當(dāng)group內(nèi)消費(fèi)者數(shù)量變化的時(shí)候不會(huì)有reblance行為發(fā)生跃须。assign的方法不能和subscribe方法同時(shí)使用站叼。
pom文件
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.yqs</groupId>
<artifactId>demo</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>demo</name>
<description>Demo project for Spring Boot</description>
<properties>
<java.version>1.8</java.version>
</properties>
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-aop</artifactId>
<version>1.5.21.RELEASE</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter</artifactId>
<version>1.5.21.RELEASE</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
<version>1.5.21.RELEASE</version>
<exclusions>
<exclusion>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-logging</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<version>1.5.21.RELEASE</version>
<scope>test</scope>
<exclusions>
<exclusion>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-logging</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-autoconfigure</artifactId>
<version>1.5.21.RELEASE</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot</artifactId>
<version>1.5.21.RELEASE</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-log4j2</artifactId>
<version>1.5.21.RELEASE</version>
</dependency>
<dependency>
<groupId>org.springframework.kafka</groupId>
<artifactId>spring-kafka</artifactId>
<version>1.3.9.RELEASE</version>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.73</version>
</dependency>
<!-- https://mvnrepository.com/artifact/log4j/log4j -->
<dependency>
<groupId>log4j</groupId>
<artifactId>log4j</artifactId>
<version>1.2.14</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.commons/commons-lang3 -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-lang3</artifactId>
<version>3.4</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<version>3.6.1</version>
<configuration>
<source>1.8</source>
<target>1.8</target>
</configuration>
</plugin>
</plugins>
</build>
</project>
kafka消費(fèi)者 :
通過(guò)switchOn變量來(lái)手動(dòng)的控制暫停跟恢復(fù)
@Component
public class KafkaPauseTest implements ApplicationListener<ContextRefreshedEvent> {
static Logger log = Logger.getLogger(KafkaPauseTest.class);
@Autowired
private Consumer consumer;
TopicPartition partition0 = new TopicPartition("test_topic", 0);
@Override
public void onApplicationEvent(ContextRefreshedEvent event) {
new Thread(()->{
try {
consumer.assign(Arrays.asList(new TopicPartition[]{partition0}));
while (true) {
log.info("獲取開(kāi)關(guān)變量="+ SwitchController.switchOn);
if (SwitchController.switchOn == 1){
log.info("暫停消費(fèi)開(kāi)始");
consumer.pause(Arrays.asList(new TopicPartition[]{partition0}));
log.info("暫停消費(fèi)結(jié)束");
}else {
log.info("恢復(fù)消費(fèi)開(kāi)始");
consumer.resume(Collections.singletonList(partition0));
log.info("恢復(fù)消費(fèi)結(jié)束");
}
ConsumerRecords<String, String> records = consumer.poll(5000);
log.info("poll結(jié)果結(jié)束");
for (ConsumerRecord<String, String> record : records) {
log.info("topic = " + record.topic() + ", partition = " + record.partition());
log.info("offset = " + record.offset());
log.info("value = " + record.value());
String msg = record.value();
}
consumer.commitSync();
}
} finally {
consumer.close();
}
}).start();
}
}
控制類:
@RestController
public class SwitchController {
public static volatile int switchOn = 0;
@RequestMapping(value = "/pause", method = RequestMethod.GET)
public String pause() {
switchOn = 1;
return "暫停監(jiān)聽(tīng)";
}
@RequestMapping(value = "/reuse", method = RequestMethod.GET)
public String reuse() {
switchOn = 0;
return "恢復(fù)監(jiān)聽(tīng)";
}
}