Kafka(二)常用脚本命令整理

Kafka(二)常用脚本命令整理

1.1 Kafka服务端常用的脚本

首先先启动Kafka服务,因为我本地是用docker-compose安装的,所以直接进去docker容器

1
2
3
4
5
docker exec -it kafka bash

# kafka的服务端脚本位置在 /opt/bitnami/kafka/bin

cd /opt/bitnami/kafka/bin
1
2
3
4
5
6
7
connect-distributed.sh	      kafka-configs.sh		   kafka-dump-log.sh	     kafka-metadata-quorum.sh	    kafka-server-start.sh		kafka-verifiable-producer.sh
connect-mirror-maker.sh kafka-console-consumer.sh kafka-e2e-latency.sh kafka-metadata-shell.sh kafka-server-stop.sh trogdor.sh
connect-plugin-path.sh kafka-console-producer.sh kafka-features.sh kafka-mirror-maker.sh kafka-storage.sh windows
connect-standalone.sh kafka-consumer-groups.sh kafka-get-offsets.sh kafka-producer-perf-test.sh kafka-streams-application-reset.sh zookeeper-security-migration.sh
kafka-acls.sh kafka-consumer-perf-test.sh kafka-jmx.sh kafka-reassign-partitions.sh kafka-topics.sh zookeeper-server-start.sh
kafka-broker-api-versions.sh kafka-delegation-tokens.sh kafka-leader-election.sh kafka-replica-verification.sh kafka-transactions.sh zookeeper-server-stop.sh
kafka-cluster.sh kafka-delete-records.sh kafka-log-dirs.sh kafka-run-class.sh kafka-verifiable-consumer.sh zookeeper-shell.sh
1
2
3
4
5
6
7
8
9
10
11
12
13
# topic主题相关操作

//创建一个新的topic
./kafka-topics.sh --bootstrap-server localhost:9092 --topic test-topic --create

//列出已经存在的topic列表
./kafka-topics.sh --bootstrap-server localhost:9092 --list

//查看指定topic主题详情
./kafka-topics.sh --bootstrap-server localhost:9092 --topic test-topic --describe

//修改指定topic信息
./kafka-topics.sh --bootstrap-server localhost:9092 --topic test-topic --alter --partitions 2
1
2
3
4
# 生产者和消费者

./kafka-console-producer.sh --bootstrap-server localhost:9092 --topic test-topic //生产者
./kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic test-topic //消费者
2.1 java操作kafka

maven依赖

1
2
3
4
5
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.6.1</version>
</dependency>
2.2 生产端代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
package com.hqd8080.test.producer;

import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

import java.util.HashMap;
import java.util.Map;

public class KafkaProducerTest {
public static void main(String[] args) {
//1.创建配置对象
//2.创建生产者对象
//3.创建数据
//4.通过生产者对象将数据发送到kafka
//5.关闭生产者对象
Map<String, Object> producerConfig = new HashMap<>();
producerConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
producerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
producerConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

KafkaProducer<String, String> producer = new KafkaProducer<String, String>(producerConfig);

for (int i =1; i <= 10; i++) {
ProducerRecord<String, String> record = new ProducerRecord<String, String>(
"test-topic","key-"+i,"this is a test-"+i
);
producer.send(record);
}

producer.close();
}
}
2.3 消费端代码
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
package com.hqd8080.test.consumer;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;

import java.util.Collections;
import java.util.HashMap;
import java.util.Map;

public class KafkaConsumerTest {
public static void main(String[] args) {
//1.创建配置对象
//2.创建消费者对象
//3.订阅主题
//4.消费者对象从kafka主题中拉取数据
//5.关闭消费者对象
Map<String, Object> consumerConfig = new HashMap<>();
consumerConfig.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
consumerConfig.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerConfig.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerConfig.put(ConsumerConfig.GROUP_ID_CONFIG, "consumer-group-id");

KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(consumerConfig);
consumer.subscribe(Collections.singleton("test-topic"));

while (true) {
final ConsumerRecords<String, String> records = consumer.poll(100);
for (ConsumerRecord<String, String> record : records) {
System.out.println(record);
}
}
//consumer.close();
}
}
2.4 代码创建Topic主题

kafka broker中的config/server.properties配置文件中配置了auto.create.topics.enable参数为true(默认值就是true)
那么当生产者向一个尚未创建的topic发送消息时,会自动创建一个num.partitions(默认值为1)个分区和default.replication.factor(默认值为1)个副本的对应topic

不过一般不建议将auto.create.topics.enable参数设置为true,因为这个参数会影响topic的管理与维护

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
package com.hqd8080.test.admin;

import org.apache.kafka.clients.admin.Admin;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.CreateTopicsResult;
import org.apache.kafka.clients.admin.NewTopic;

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;

public class AdminTopicTest {
public static void main(String[] args) {

Map<String, Object> confMap = new HashMap<>();
confMap.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");

final Admin admin = Admin.create(confMap);

String topicName = "hqd8080-topic"; //主题、逻辑分区
int partitionCount = 1; //分区
short replicationCount = 1; //副本

NewTopic topic = new NewTopic(topicName, partitionCount, replicationCount);

final CreateTopicsResult result = admin.createTopics(
Arrays.asList(topic)
);

admin.close();
}
}