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) { 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) { 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); } } } }
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(); } }