预见猿份
主题
首页面试题模拟面试SQL练习在线工具关于我们老苗一对一私教学员评价
实战项目
零基础学习路线项目前置基础创新WMS项目Java微服务框架与实战云岚到家项目闪聚支付项目学成在线项目青橙电商项目JVM原理与实战调优分布式事务专题Java高频面试题MySQL从入门到精通Java数据结构与算法Java 并发编程(JUC)实战课程老苗一对一私教学员评价
blog我的会员
  • Kafka 分布式消息中间件实战

    • 课程介绍
    • 第一章 初识 Kafka
    • 第二章 Producer 生产者详解
    • 第三章 Consumer 消费者详解
    • 第四章 Topic 主题管理
    • 第五章 Consumer 消费者详解
    • 第六章 Kafka 存储原理
    • 第七章 Kafka 稳定性与可靠性
    • 第八章 Kafka 高级应用
    • 第九章 Kafka 集群与监控
    • 第十章 配套面试专题







----- 到底线了 -----

Kafka从入门到应用 ​

第4章 主题 ​

tips 学完这一章你可以

  • 深入学习Kafka主题的管理
  • KafkaAdminClient应用

4.1 管理 ​

4.1.1 创建主题 ​

命令

bin/kafka-topics.sh --zookeeper localhost:2181 --create --topic heima --partitions 2 --replicationfactor 1

localhost:2181 zookeeper所在的ip,zookeeper 必传参数,多个zookeeper用 ‘,’分开。

partitions 用于设置主题分区数,每个线程处理一个分区数据

replication-factor 用于设置主题副本数,每个副本分布在不通节点,不能超过总结点数。如你只有一个节点,但是创建时指定副本数为2,就会报错。

查看topic元数据信细的方法

topic元数据信息保存在Zookeeper节点中

bash
// 连接zk client
itcast@Server-node:/mnt/d/zookeeper-3.4.14$ bin/zkCli.sh -server localhost:2181
Connecting to localhost:2181
...........................................
[zk: localhost:2181(CONNECTED) 2] get /brokers/topics/heima
{"version":1,"partitions":{"1":[0],"0":[0]}}
cZxid = 0x618
ctime = Wed Aug 28 05:51:35 GMT 2019
mZxid = 0x618
mtime = Wed Aug 28 05:51:35 GMT 2019
pZxid = 0x619
cversion = 1
dataVersion = 0
aclVersion = 0
ephemeralOwner = 0x0
dataLength = 44
numChildren = 1
[zk: localhost:2181(CONNECTED) 3]
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18

4.1.2 查看主题 ​

bash
//  查看所有主题
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --list --
zookeeper localhost:2181
__consumer_offsets
__transaction_state
_schemas
heima
topic0701
topic0703
topic0703kafka_source
topic0828
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$
1
2
3
4
5
6
7
8
9
10
11
12
bash
//  查看某个特定主题信息,不指定topic则查询所有  通过 --describe
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --describe --
zookeeper localhost:2181 --topic heima
Topic:heima     PartitionCount:2        ReplicationFactor:1     Configs:
Topic: heima    Partition: 0    Leader: 0       Replicas: 0     Isr: 0
Topic: heima    Partition: 1    Leader: 0       Replicas: 0     Isr: 0
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$
//查看正在同步的主题
// 通过 --describe 和 under-replicated-partitions命令组合查看 under-replacation状态
1
2
3
4
5
6
7
8
9

4.1.3 修改主题 ​

bash
// 增加配置
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --alter --
zookeeper localhost:2181 --topic heima
--config flush.messages=1
WARNING: Altering topic configuration from this script has been deprecated and
may be removed in future releases.
Going forward, please use kafka-configs.sh for this functionality
Updated config for topic heima.
1
2
3
4
5
6
7
8
bash
// 删除配置
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --alter --
zookeeper localhost:2181 --topic heima --delete-config flush.messages
WARNING: Altering topic configuration from this script has been deprecated and
may be removed in future releases.
Going forward, please use kafka-configs.sh for this functionality
Updated config for topic heima.
1
2
3
4
5
6
7

4.1.4 删除主题 ​

若 delete.topic.enable=true 直接彻底删除该 Topic。 若 delete.topic.enable=false 如果当前Topic 没有使用过即没有传输过信息:可以彻底删除。 如果当前 Topic 有使用过即有过传输过信息:并没有真正删除 Topic 只是把这个 Topic 标记为删除(marked for deletion),重启 Kafka Server 后删除。

bash
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --delete --
zookeeper localhost:2181 --topic heima
Topic heima is marked for deletion.
Note: This will have no impact if delete.topic.enable is not set to true.
1
2
3
4
bash
// 标记为 marked for deletion
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --list --
zookeeper localhost:2181
__consumer_offsets
topic0701
heima - marked for deletion
1
2
3
4
5
6

4.2 增加分区 ​

bash
// 增加分区数
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --alter --
zookeeper localhost:2181 --topic heima --partitions 3
WARNING: If partitions are increased for a topic that has a key, the partition
logic or ordering of the messages will be affected
Adding partitions succeeded!
1
2
3
4
5
6
bash
//修改分区数时,仅能增加分区个数。若是用其减少 partition 个数,则会报如下错误信息:
itcast@Server-node:/mnt/d/kafka_2.12-2.2.1$ bin/kafka-topics.sh --alter --
zookeeper localhost:2181 --topic heima --partitions 2
WARNING: If partitions are increased for a topic that has a key, the partition
logic or ordering of the messages will be affected
Error while executing topic command : The number of partitions for a topic can
only be increased. Topic heima currently has 3 partitions, 2 would not be an
increase.
[2019-08-28 08:43:41,478] ERROR
org.apache.kafka.common.errors.InvalidPartitionsException: The number of
partitions for a topic can only be increased. Topic heima currently has 3
partitions, 2 would not be an increase.
(kafka.admin.TopicCommand$)
1
2
3
4
5
6
7
8
9
10
11
12
13

4.3 分区副本的分配--只做了解 ​

4.4 其他主题参数配置 ​

见官方文档:http://kafka.apache.org/documentation/#topicconfigs

text
Configurations pertinent to topics have both a server default as well an
optional per-topic override. If no per-topic configuration is given the server
default is used. The override can be set at topic creation time by giving one or
more --config options. This example creates a topic named my-topic with a custom
max message size and flush rate:
1
2
> bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic my-
topic --partitions 1 \
--replication-factor 1 --config max.message.bytes=64000 --config
flush.messages=1
Overrides can also be changed or set later using the alter configs command. This
example updates the max message size for my-topic:
1
2
> bin/kafka-configs.sh --zookeeper localhost:2181 --entity-type topics --entity-
name my-topic
--alter --add-config max.message.bytes=128000
To check overrides set on the topic you can do
1
> bin/kafka-configs.sh --zookeeper localhost:2181 --entity-type topics --entity-
name my-topic --describe
To remove an override you can do
1
2
> bin/kafka-configs.sh --zookeeper localhost:2181  --entity-type topics --
entity-name my-topic
--alter --delete-config max.message.bytes
The following are the topic-level configurations. The server's default
configuration for this property is given under the Server Default Property
heading. A given server default config value only applies to a topic if it does
not have an explicit topic config override.
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

4.5 KafkaAdminClient应用 ​

我们都习惯使用Kafka中bin目录下的脚本工具来管理查看Kafka,但是有些时候需要将某些管理查看的功能集成到系统(比如Kafka Manager)中,那么就需要调用一些API来直接操作Kafka了。

见代码库:com.yjoffer.kafka.chapter4.KafkaAdminConfigOperation

text
/**
* KafkaAdminClient应用
*/
public class KafkaAdminConfigOperation {
1
2
3
4
text
public static void main(String[] args) throws ExecutionException,
InterruptedException {
//        describeTopicConfig();
//        alterTopicConfig();
addTopicPartitions();
1
2
3
4
5
text
}
1
text
//Config(entries=[ConfigEntry(name=compression.type, value=producer,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=leader.replication.throttled.replicas, value=,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=message.downconversion.enable, value=true,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=min.insync.replicas, value=1, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=segment.jitter.ms, value=0, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=cleanup.policy, value=delete, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]), ConfigEntry(name=flush.ms,
value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false,
isReadOnly=false, synonyms=[]),
ConfigEntry(name=follower.replication.throttled.replicas, value=,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=segment.bytes, value=1073741824, source=STATIC_BROKER_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=retention.ms, value=604800000, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=flush.messages, value=9223372036854775807,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=message.format.version, value=2.0-IV1, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=file.delete.delay.ms, value=60000, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=max.message.bytes, value=1000012, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=min.compaction.lag.ms, value=0, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=message.timestamp.type, value=CreateTime,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=preallocate, value=false, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=min.cleanable.dirty.ratio, value=0.5, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=index.interval.bytes, value=4096, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=unclean.leader.election.enable, value=false,
source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=retention.bytes, value=-1, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]),
ConfigEntry(name=delete.retention.ms, value=86400000, source=DEFAULT_CONFIG,
isSensitive=false, isReadOnly=false, synonyms=[]), ConfigEntry(name=segment.ms,
value=604800000, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false,
synonyms=[]), ConfigEntry(name=message.timestamp.difference.max.ms,
value=9223372036854775807, source=DEFAULT_CONFIG, isSensitive=false,
isReadOnly=false, synonyms=[]), ConfigEntry(name=segment.index.bytes,
value=10485760, source=DEFAULT_CONFIG, isSensitive=false, isReadOnly=false,
synonyms=[])])
public static void describeTopicConfig() throws ExecutionException,
InterruptedException {
String brokerList =  "localhost:9092";
String topic = "heima";
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
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
text
Properties props = new Properties();
1
text
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000);
AdminClient client = AdminClient.create(props);
1
2
3
text
ConfigResource resource =
new ConfigResource(ConfigResource.Type.TOPIC, topic);
DescribeConfigsResult result =
client.describeConfigs(Collections.singleton(resource));
Config config = result.all().get().get(resource);
System.out.println(config);
client.close();
}
1
2
3
4
5
6
7
8
text
public static void alterTopicConfig() throws ExecutionException,
InterruptedException {
String brokerList =  "localhost:9092";
String topic = "heima";
1
2
3
4
text
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000);
AdminClient client = AdminClient.create(props);
1
2
3
4
text
ConfigResource resource =
new ConfigResource(ConfigResource.Type.TOPIC, topic);
ConfigEntry entry = new ConfigEntry("cleanup.policy", "compact");
Config config = new Config(Collections.singleton(entry));
Map<ConfigResource, Config> configs = new HashMap<>();
configs.put(resource, config);
AlterConfigsResult result = client.alterConfigs(configs);
result.all().get();
1
2
3
4
5
6
7
8
text
client.close();
}
1
2
text
public static void addTopicPartitions() throws ExecutionException,
InterruptedException {
String brokerList =  "localhost:9092";
String topic = "heima";
1
2
3
4
text
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList);
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000);
AdminClient client = AdminClient.create(props);
1
2
3
4
text
NewPartitions newPartitions = NewPartitions.increaseTo(5);
Map<String, NewPartitions> newPartitionsMap = new HashMap<>();
newPartitionsMap.put(topic, newPartitions);
CreatePartitionsResult result =
client.createPartitions(newPartitionsMap);
result.all().get();
1
2
3
4
5
6
text
client.close();
}
}
1
2
3

总结 ​

本章主要讲解了Kafka两大核心之一:主题,通过对主题的增删该查、配置等内容来了解主题的相关内容。

← 第三章 Consumer 消费者详解第五章 Consumer 消费者详解 →








如果发现文档内容有错误或排版错乱,请及时联系站长老苗修改,不胜感激。联系我们
关于我们 | 隐私政策 | 豫ICP备2026003386号-4 | 豫公网安备41010202004008号
目录

本页无章节