增加kafka 测试

This commit is contained in:
hz21056617
2022-04-15 11:08:06 +08:00
parent 2640a2c61c
commit 33a58e8aac
11 changed files with 671 additions and 0 deletions
@@ -4,6 +4,7 @@ import lombok.AllArgsConstructor;
import lombok.SneakyThrows;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ConcurrentLinkedQueue;
@AllArgsConstructor
public class BlockingQueueTestWork extends Thread {
@@ -27,6 +28,7 @@ public class BlockingQueueTestWork extends Thread {
}
public void offerTest() {
ConcurrentLinkedQueue<Integer> concurrentLinkedQueue = new ConcurrentLinkedQueue<>();
boolean offer = arrayBlockingQueue.offer(1);
if (!offer) {
System.out.println("队列满了");
+23
View File
@@ -0,0 +1,23 @@
<?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 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>JavaBasiceDemo</artifactId>
<groupId>com.wyl.example</groupId>
<version>1.0-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>kafka</artifactId>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
</dependency>
</dependencies>
</project>
@@ -0,0 +1,51 @@
package com.wyl.kafka;
/**
* kafka配置
* @ClassName: KafkaConfig
* @Date: 2022/4/12 14:11
* @author wangyl
* @version V1.0
*/
public class KafkaConfigConstant {
/**
*kafka的broker地址
*/
public static final String BOOTSTRAP_SERVERS = "localhost:9092";
/**
*kafka的topic
*/
public static final String TOPIC = "test";
/**
*kafka的groupId
*/
public static final String GROUP_ID = "test-group";
/**
*kafka的key的序列化方式
*/
public static final String KEY_DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer";
/**
*kafka的value的序列化方式
*/
public static final String VALUE_DESERIALIZER = "org.apache.kafka.common.serialization.StringDeserializer";
/**
* ACK策略
*/
public static final String ACKS = "all";
/**
*kafka的消息发送最大重试次数
*/
public static final String RETRIES = "3";
/**
*kafka的消息发送最大延时
*/
public static final String BATCH_SIZE = "16384";
/**
*kafka的消息发送最大延时
*/
public static final String LINGER_MS = "1";
public static final String BUFFER_MEMORY = "33554432";
public static final String AUTO_COMMIT_INTERVAL = "1000";
public static final String AUTO_OFFSET_RESET = "earliest";
public static final String ENABLE_AUTO_COMMIT = "true";
}
@@ -0,0 +1,50 @@
package com.wyl.kafka.consumer;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import java.util.Properties;
/**
* kafak生产者示例
* @ClassName: KafkaProducer
* @Date: 2022/4/12 14:36
* @author wangyl
* @version V1.0
*/
public class KafkaConsumerExample {
/**
* 通用消费者
*/
public static Consumer<String, String> KafkaConsumerNormal() {
Properties props = NormalProperties();
Consumer<String, String> consumerNormal = new KafkaConsumer<>(props);
return consumerNormal;
}
/**
* 手动提交消费者
*/
public static Consumer<String, String> KafkaConsumerManualCommit() {
Properties props = NormalProperties();
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
Consumer<String, String> consumerNormal = new KafkaConsumer<>(props);
return consumerNormal;
}
/**
* 通用配置
*/
private static Properties NormalProperties() {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "172.20.0.187:9092,172.20.0.188:9092,172.20.0.189:9092");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "wylll");
return props;
}
}
@@ -0,0 +1,222 @@
package com.wyl.kafka.consumer;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import java.time.Duration;
import java.util.*;
import java.util.stream.Stream;
/**
* kafka发送消息测试
* @ClassName: KafkaProducerSend
* @Date: 2022/4/12 15:18
* @author wangyl
* @version V1.0
*/
@Slf4j
public class KafkaConsumerPoll {
/**
* 消费者消费消息
* @param kafkaConsumerNormal
* @param topic
* @return void
* @Date 2022/4/12 17:43
* @Author wangyl
* @Version V1.0
*/
public static void PollMessage(Consumer<String, String> kafkaConsumerNormal, String topic) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic));
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
}
}
/**
* 同步手动提交
* @param kafkaConsumerNormal
* @param topic
* @return void
* @Date 2022/4/13 18:44
* @Author wangyl
* @Version V1.0
*/
public static void PollMessageCommitSync(Consumer<String, String> kafkaConsumerNormal, String topic) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic));
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
kafkaConsumerNormal.commitSync();
}
}
/**
* 同步手动提交offset
* @param kafkaConsumerNormal
* @param topic
* @return void
* @Date 2022/4/13 18:44
* @Author wangyl
* @Version V1.0
*/
public static void PollMessageCommitSyncOffsets(Consumer<String, String> kafkaConsumerNormal, String topic) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic));
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1));
});
kafkaConsumerNormal.commitSync(offsets);
offsets.clear();
}
}
/**
* 异步提交,有回调
* @param kafkaConsumerNormal
* @param topic
* @return void
* @Date 2022/4/13 18:47
* @Author wangyl
* @Version V1.0
*/
public static void PollMessageCommitAsyncCallBack(Consumer<String, String> kafkaConsumerNormal, String topic) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic));
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
offsets.put(new TopicPartition(record.topic(), record.partition()), new OffsetAndMetadata(record.offset() + 1));
});
kafkaConsumerNormal.commitAsync(offsets, (offsetsInfo, exception) -> {
if (exception != null) {
log.error("提交消息异常:{}", exception.getMessage());
} else {
log.info("提交消息成功:{}", offsetsInfo);
}
});
}
}
/**
* 异步回调offsets
* @param kafkaConsumerNormal
* @param topic
* @return void
* @Date 2022/4/13 18:50
* @Author wangyl
* @Version V1.0
*/
public static void PollMessageCommitAsyncOffsetsCallBack(Consumer<String, String> kafkaConsumerNormal, String topic) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic));
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
kafkaConsumerNormal.commitAsync();
}
}
public static void PollMessageSeek(Consumer<String, String> kafkaConsumerNormal, String topic,int partition,int offset) {
//消费者消费消息
kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
kafkaConsumerNormal.seek(new TopicPartition(topic, partition), offset);
}
});
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
kafkaConsumerNormal.commitAsync();
}
}
public static void PollMessageSeekBegin(Consumer<String, String> kafkaConsumerNormal, String topic,int partition) {
final TopicPartition topicPartition = new TopicPartition(topic, partition);
final List<TopicPartition> topicPartitions = Arrays.asList(topicPartition);
kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
kafkaConsumerNormal.seekToBeginning(topicPartitions);
}
});
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
kafkaConsumerNormal.commitAsync();
}
}
public static void PollMessageSeekEnd(Consumer<String, String> kafkaConsumerNormal, String topic,int partition) {
//消费者消费消息
final TopicPartition topicPartition = new TopicPartition(topic, partition);
final List<TopicPartition> topicPartitions = Arrays.asList(topicPartition);
kafkaConsumerNormal.subscribe(Arrays.asList(topic), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
kafkaConsumerNormal.seekToEnd(topicPartitions);
}
});
while (true) {
//消费消息
kafkaConsumerNormal.poll(Duration.ofMillis(100))
.forEach(record -> {
log.info("消费者消费消息:{},分区:{},topic{},offset:{}", record.value(), record.partition(), record.topic(), record.offset());
});
kafkaConsumerNormal.commitAsync();
}
}
}
@@ -0,0 +1,55 @@
package com.wyl.kafka.producers;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import java.util.Properties;
/**
* kafak生产者示例
* @ClassName: KafkaProducer
* @Date: 2022/4/12 14:36
* @author wangyl
* @version V1.0
*/
public class KafkaProducerExample {
/**
* 通用生产者
*/
public static Producer<String, String> KafkaProducerNormal() {
Properties props = NormalProperties();
Producer<String, String> producerNormal = new KafkaProducer(props);
return producerNormal;
}
/**
* 事务消息生产者
*/
public static Producer<String, String> KafkaProducerTransactional() {
Properties props = NormalProperties();
/**
*设置事务 id(必须),事务 id 任意起名
*/
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "transactionalId");
props.put(ProducerConfig.RETRIES_CONFIG, 2);
Producer<String, String> producerNormal = new KafkaProducer(props);
return producerNormal;
}
/**
* 通用配置
*/
private static Properties NormalProperties() {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "172.20.0.187:9092,172.20.0.188:9092,172.20.0.189:9092");
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
props.put(ProducerConfig.LINGER_MS_CONFIG, 1);
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
return props;
}
}
@@ -0,0 +1,79 @@
package com.wyl.kafka.producers;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
/**
* kafka发送消息测试
* @ClassName: KafkaProducerSend
* @Date: 2022/4/12 15:18
* @author wangyl
* @version V1.0
*/
@Slf4j
public class KafkaProducerSend {
/**
* 生产者发送回调消息
* @param kafkaProducerNormal
* @param topic
* @param message
* @return void
* @Date 2022/4/12 17:43
* @Author wangyl
* @Version V1.0
*/
public static void sendMessageAndCallback(Producer<String, String> kafkaProducerNormal, String topic, String message) {
ProducerRecord<String, String> messageObj = new ProducerRecord<>(topic, message);
kafkaProducerNormal.send(messageObj, (metadata, exception) -> {
if (null != exception) {
log.error("kafka 信息推送失败[" + messageObj + "]", exception);
}
log.info("消息发送成功[{}]", messageObj);
});
kafkaProducerNormal.close();
}
/**
* 发送事务消息
* @param kafkaProducerTransactional
* @param topic
* @param message
* @return void
* @Date 2022/4/12 17:42
* @Author wangyl
* @Version V1.0
*/
public static void sendMessageTransactional(Producer<String, String> kafkaProducerTransactional, String topic, String message) {
/**
* 初始化事务
*/
kafkaProducerTransactional.initTransactions();
/**
* 开启事务
*/
kafkaProducerTransactional.beginTransaction();
log.info("开始发送事务消息");
try {
for (int i = 0; i < 10; i++){
ProducerRecord<String, String> messageObj = new ProducerRecord<>(topic, message + i);
kafkaProducerTransactional.send(messageObj, (metadata, exception) -> {
if (null != exception) {
log.error("kafka 信息推送失败[" + messageObj + "]", exception);
}
});
}
log.info("消息发送成功");
kafkaProducerTransactional.commitTransaction();
}
catch (Exception e) {
log.error("kafka 事务消息推送失败", e);
kafkaProducerTransactional.abortTransaction();
}
finally {
kafkaProducerTransactional.close();
}
}
}
+59
View File
@@ -0,0 +1,59 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration scan="true" scanPeriod="120 seconds" debug="false">
<include resource="org/springframework/boot/logging/logback/defaults.xml"/>
<property name="logback.output" value="./logs"/>
<property name="root_path" value="./logs"/>
<appender name="console" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%date{yyyy-MM-dd HH:mm:ss.SSS} java %-5level [%replace(%thread){'\s', '-'}] %logger{36} %X{trace.id} %msg%n
</pattern>
<charset>utf8</charset>
</encoder>
</appender>
<appender name="async-console" class="ch.qos.logback.classic.AsyncAppender">
<appender-ref ref="console"/>
<includeCallerData>true</includeCallerData>
</appender>
<appender name="rollingFile" class="ch.qos.logback.core.rolling.RollingFileAppender">
<file>${logback.output}/server.log</file>
<rollingPolicy class="ch.qos.logback.core.rolling.SizeAndTimeBasedRollingPolicy">
<fileNamePattern>${logback.output}/server.%d{yyyy-MM-dd}.%i.log
</fileNamePattern>
<maxFileSize>1GB</maxFileSize>
<maxHistory>10</maxHistory>
<totalSizeCap>10GB</totalSizeCap>
</rollingPolicy>
<encoder>
<pattern>%date{HH:mm:ss.SSS} [%replace(%thread){'\s', '-'}] %-5level %logger{36} %X{trace.id:-huizTraceId} %msg%n
</pattern>
</encoder>
</appender>
<appender name="async-rollingFile" class="ch.qos.logback.classic.AsyncAppender">
<appender-ref ref="rollingFile"/>
<includeCallerData>true</includeCallerData>
</appender>
<logger name="java.sql.Connection" level="DEBUG"/>
<logger name="java.sql.Statement" level="DEBUG"/>
<logger name="java.sql.PreparedStatement" level="DEBUG"/>
<logger name="java.sql.ResultSet" level="OFF"/>
<logger name="jdbc.sqlonly" level="OFF"/>
<logger name="jdbc.audit" level="OFF"/>
<logger name="jdbc.connection" level="OFF"/>
<logger name="jdbc.resultset" level="OFF"/>
<logger name="com.hzins.bsp" level="OFF"/>
<logger name="com.hzins.remoting" level="OFF"/>
<logger name="com.hzins.remoting" level="OFF"/>
<logger name="org.apache.http" level="OFF"/>
<root level="info">
<appender-ref ref="async-console"/>
<appender-ref ref="async-rollingFile"/>
</root>
</configuration>
@@ -0,0 +1,88 @@
package com.wyl.kafka.test;
import com.wyl.kafka.consumer.KafkaConsumerExample;
import com.wyl.kafka.consumer.KafkaConsumerPoll;
import com.wyl.kafka.producers.KafkaProducerExample;
import com.wyl.kafka.producers.KafkaProducerSend;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.producer.Producer;
import org.junit.jupiter.api.Test;
@Slf4j
public class ConsumerTest {
final static String TOPIC_NAME = "wyl-filebeat-test";
/**
* 消费者模式
*/
@Test
public void normalProducerSendMessageAndCallbackTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerNormal();
KafkaConsumerPoll.PollMessage(consumer, TOPIC_NAME);
}
/**
* 手动同步提交offset
*/
@Test
public void pollMessageCommitSyncTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerManualCommit();
KafkaConsumerPoll.PollMessageCommitSync(consumer, TOPIC_NAME);
}
/**
* 手动异步提交offset
*/
@Test
public void pollMessageCommitSyncOffsets() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerManualCommit();
KafkaConsumerPoll.PollMessageCommitSyncOffsets(consumer, TOPIC_NAME);
}
/**
* 手动异步提交offset回调
*/
@Test
public void pollMessageCommitAsyncCallBackTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerManualCommit();
KafkaConsumerPoll.PollMessageCommitAsyncCallBack(consumer, TOPIC_NAME);
}
/**
* 手动异步提交offset批量回调
*/
@Test
public void pollMessageCommitAsyncOffsetsCallBackTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerManualCommit();
KafkaConsumerPoll.PollMessageCommitAsyncOffsetsCallBack(consumer, TOPIC_NAME);
}
/**
* 移动到指定问位置消费消息
*/
@Test
public void pollMessageSeekTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerNormal();
KafkaConsumerPoll.PollMessageSeek(consumer, TOPIC_NAME, 1, 34236);
}
/**
* 移动到开始位置消费消息
*/
@Test
public void pollMessageSeekBeginTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerNormal();
KafkaConsumerPoll.PollMessageSeekBegin(consumer, TOPIC_NAME, 1);
}
/**
* 移动到结束位置消费消息
*/
@Test
public void pollMessageCommitSeekEndTest() {
final Consumer<String, String> consumer = KafkaConsumerExample.KafkaConsumerNormal();
KafkaConsumerPoll.PollMessageSeekEnd(consumer, TOPIC_NAME, 1);
}
}
@@ -0,0 +1,24 @@
package com.wyl.kafka.test;
import com.wyl.kafka.producers.KafkaProducerExample;
import com.wyl.kafka.producers.KafkaProducerSend;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.producer.Producer;
import org.junit.jupiter.api.Test;
@Slf4j
public class ProducerTest {
final static String TOPIC_NAME = "wyl-filebeat-test";
@Test
public void normalProducerSendMessageAndCallbackTest() {
final Producer<String, String> stringStringProducer = KafkaProducerExample.KafkaProducerNormal();
KafkaProducerSend.sendMessageAndCallback(stringStringProducer, TOPIC_NAME, "test123");
}
@Test
public void transactionalProducerSendMessageAndCallbackTest() {
final Producer<String, String> stringStringProducer = KafkaProducerExample.KafkaProducerTransactional();
KafkaProducerSend.sendMessageTransactional(stringStringProducer, TOPIC_NAME, "test123");
}
}
+18
View File
@@ -12,6 +12,7 @@
<module>elasticsearch</module>
<module>concurrent</module>
<module>file-read</module>
<module>kafka</module>
</modules>
<packaging>pom</packaging>
<properties>
@@ -25,6 +26,12 @@
<artifactId>cglib</artifactId>
<version>3.3.0</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.apache.kafka/kafka-clients -->
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>2.8.1</version>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
@@ -39,5 +46,16 @@
<version>5.8.2</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-core</artifactId>
<version>1.2.3</version>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
<version>1.2.3</version>
</dependency>
</dependencies>
</project>