📌 本文是「Kafka 3.9 教程」系列第 九 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
- 9.1 Producer API
- 9.2 Consumer API
- 9.3 Streams API
- 9.4 Connect API
- 9.5 Admin API
- 9.6 交互式查询
- 9.7 应用重置工具
9.1 Producer API
9.1.1 概述
Producer API 用于将消息发布到 Kafka 主题。它是 Kafka 最基础的 API 之一。
9.1.2 Maven 依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.9.0</version>
</dependency>
9.1.3 ProducerRecord
ProducerRecord 是发送到 Kafka 的基本消息单元。
构造方法
// 1. 仅指定主题和值
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", "value");
// 2. 指定主题、键和值
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", "key1", "value");
// 3. 指定主题、分区、键和值
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", 0, "key1", "value");
// 4. 指定主题、分区、时间戳、键和值
ProducerRecord<String, String> record = new ProducerRecord<>(
"my-topic", 0, System.currentTimeMillis(), "key1", "value"
);
// 5. 指定主题、分区、时间戳、键、值和头部
Headers headers = new RecordHeaders();
headers.add("header-key", "header-value".getBytes());
ProducerRecord<String, String> record = new ProducerRecord<>(
"my-topic", 0, System.currentTimeMillis(), "key1", "value", headers
);
常用方法
| 方法 | 描述 |
|---|---|
topic() | 获取主题名称 |
partition() | 获取分区(可能为 null) |
key() | 获取消息键 |
value() | 获取消息值 |
timestamp() | 获取时间戳 |
headers() | 获取消息头 |
9.1.4 KafkaProducer
KafkaProducer 是线程安全的,可以在多个线程间共享。
创建 Producer
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
// 可选配置
props.put( , "all");
props.put(ProducerConfig.RETRIES_CONFIG, 3);
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
发送消息
同步发送:
try {
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", "key", "value");
RecordMetadata metadata = producer.send(record).get();
System.out.println("消息发送成功:");
System.out.println("主题: " + metadata.topic());
System.out.println("分区: " + metadata.partition());
System.out.println("偏移量: " + metadata.offset());
} catch (InterruptedException | ExecutionException e) {
e.printStackTrace();
} finally {
producer.close();
}
异步发送:
ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", "key", "value");
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception == null) {
System.out.println("消息发送成功: " + metadata.offset());
} else {
System.err.println("消息发送失败: " + exception.getMessage());
}
}
});
// Lambda 风格
producer.send(record, (metadata, exception) -> {
if (exception == null) {
System.out.println("发送成功: partition=" + metadata.partition() +
", offset=" + metadata.offset());
} else {
exception.printStackTrace();
}
});
事务性发送
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "my-transactional-id");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// 初始化事务
producer.initTransactions();
try {
// 开启事务
producer.beginTransaction();
// 发送消息
producer.send(new ProducerRecord<>("topic1", "key1", "value1"));
producer.send(new ProducerRecord<>("topic2", "key2", "value2"));
// 提交事务
producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException |
AuthorizationException e) {
// 致命错误,关闭 producer
producer.close();
} catch (KafkaException e) {
// 中止事务
producer.abortTransaction();
} finally {
producer.close();
}
刷新消息
// 立即发送所有缓冲的消息
producer.flush();
关闭 Producer
// 优雅关闭,等待所有消息发送完成
producer.close();
// 带超时的关闭
producer.close(Duration.ofSeconds(10));
9.1.5 RecordMetadata
RecordMetadata 包含已发送消息的元数据。
producer.send(record, (metadata, exception) -> {
if (exception == null) {
long offset = metadata.offset();
int partition = metadata.partition();
String topic = metadata.topic();
long timestamp = metadata.timestamp();
// 检查是否启用时间戳
boolean hasTimestamp = metadata.hasTimestamp();
// 序列化后的键和值的大小
int serializedKeySize = metadata.serializedKeySize();
int serializedValueSize = metadata.serializedValueSize();
}
});
9.1.6 Producer 拦截器
public class MyProducerInterceptor implements ProducerInterceptor<String, String> {
@Override
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
// 在消息发送前修改记录
return new ProducerRecord<>(
record.topic(),
record.partition(),
record.timestamp(),
record.key(),
"PREFIX: " + record.value(),
record.headers()
);
}
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
// 在收到服务端响应后调用
if (exception != null) {
System.err.println("发送失败: " + exception.getMessage());
}
}
@Override
public void close() {
// 关闭资源
}
@Override
public void configure(Map<String, ?> configs) {
// 配置拦截器
}
}
配置拦截器:
# 使用自定义拦截器
producer.interceptor.classes=com.example.MyProducerInterceptor
9.2 Consumer API
9.2.1 概述
Consumer API 用于从 Kafka 主题消费消息。
9.2.2 KafkaConsumer
KafkaConsumer 不是线程安全的,每个线程应该有自己的 consumer 实例。
创建 Consumer
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-consumer-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// 可选配置
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000);
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
订阅主题
// 1. 订阅单个主题
consumer.subscribe(Collections.singletonList("my-topic"));
// 2. 订阅多个主题
consumer.subscribe(Arrays.asList("topic1", "topic2", "topic3"));
// 3. 使用正则表达式订阅
consumer.subscribe(Pattern.compile("topic-.*"));
// 4. 使用 ConsumerRebalanceListener
consumer.subscribe(Arrays.asList("topic1", "topic2"), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
// 分区被收回前提交偏移量
consumer.commitSync();
System.out.println("分区被收回: " + partitions);
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
// 分区分配后
System.out.println("分区被分配: " + partitions);
}
});
手动分配分区
// 不使用消费者组,手动分配分区
List<TopicPartition> partitions = Arrays.asList(
new TopicPartition("my-topic", 0),
new TopicPartition("my-topic", 1)
);
consumer.assign(partitions);
// 查找偏移量
Map<TopicPartition, Long> beginningOffsets = consumer.beginningOffsets(partitions);
Map<TopicPartition, Long> endOffsets = consumer.endOffsets(partitions);
// 跳转到特定偏移量
TopicPartition partition0 = new TopicPartition("my-topic", 0);
consumer.seek(partition0, 100);
// 跳转到时间戳位置
Map<TopicPartition, Long> timestampToSearch = new HashMap<>();
timestampToSearch.put(partition0, System.currentTimeMillis() - 3600000); // 1小时前
Map<TopicPartition, OffsetAndTimestamp> offsetsForTimes = consumer.offsetsForTimes(timestampToSearch);
消费消息
简单轮询:
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("主题: %s, 分区: %d, 偏移量: %d, 键: %s, 值: %s%n",
record.topic(),
record.partition(),
record.offset(),
record.key(),
record.value()
);
}
}
} finally {
consumer.close();
}
按主题处理:
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (TopicPartition partition : records.partitions()) {
List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);
for (ConsumerRecord<String, String> record : partitionRecords) {
System.out.println("分区 " + partition + ": " + record.value());
}
long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
// 手动提交偏移量
consumer.commitSync(Collections.singletonMap(
partition, new OffsetAndMetadata(lastOffset + 1)
));
}
}
手动提交偏移量:
// 同步提交
try {
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
// 处理记录
}
// 同步提交当前偏移量
consumer.commitSync();
}
} catch (CommitFailedException e) {
e.printStackTrace();
} finally {
consumer.close();
}
// 异步提交
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
// 处理记录
}
// 异步提交
consumer.commitAsync(new OffsetCommitCallback() {
@Override
public void onComplete(Map<TopicPartition, OffsetAndMetadata> offsets,
Exception exception) {
if (exception != null) {
System.err.println("提交失败: " + exception);
}
}
});
}
ConsumerRecord
ConsumerRecord<String, String> record = ...;
// 基本信息
String topic = record.topic();
int partition = record.partition();
long offset = record.offset();
String key = record.key();
String value = record.value();
// 时间戳
long timestamp = record.timestamp();
TimestampType timestampType = record.timestampType();
// 头部
Headers headers = record.headers();
for (Header header : headers) {
System.out.println(header.key() + ": " + new String(header.value()));
}
// 序列化大小
int serializedKeySize = record.serializedKeySize();
int serializedValueSize = record.serializedValueSize();
// 检查是否有控制记录(内部使用)
boolean isControlRecord = record.isControlRecord();
消费者组管理
// 获取消费者组ID
String groupId = consumer.groupMetadata().groupId();
// 获取成员ID
String memberId = consumer.groupMetadata().memberId();
// 列出订阅的主题
Set<String> subscribed = consumer.subscription();
// 获取分配的分区
Set<TopicPartition> assignment = consumer.assignment();
// 获取消费者位置(已消费的偏移量 + 1)
Map<TopicPartition, Long> positions = new HashMap<>();
for (TopicPartition partition : assignment) {
long position = consumer.position(partition);
positions.put(partition, position);
}
// 获取提交的偏移量
Map<TopicPartition, OffsetAndMetadata> committed = consumer.committed(assignment);
// 暂停和恢复消费
consumer.pause(assignment);
// ... 做一些其他工作 ...
consumer.resume(assignment);
9.2.3 Consumer 拦截器
public class MyConsumerInterceptor implements ConsumerInterceptor<String, String> {
@Override
public ConsumerRecords<String, String> onConsume(ConsumerRecords<String, String> records) {
// 在消息返回给应用前处理
// 可以过滤、修改消息
return records;
}
@Override
public void onCommit(Map<TopicPartition, OffsetAndMetadata> offsets) {
// 在偏移量提交后调用
System.out.println("提交的偏移量: " + offsets);
}
@Override
public void close() {
// 关闭资源
}
@Override
public void configure(Map<String, ?> configs) {
// 配置拦截器
}
}
配置拦截器:
consumer.interceptor.classes=com.example.MyConsumerInterceptor
9.3 Streams API
9.3.1 概述
Kafka Streams 是用于构建流处理应用的客户端库,无需依赖外部处理系统。
9.3.2 Maven 依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-streams</artifactId>
<version>3.9.0</version>
</dependency>
9.3.3 KStream
KStream 是从主题抽象出的消息流,每条记录都是独立的。
创建 KStream
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-streams-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
// 从主题创建 KStream
KStream<String, String> source = builder.stream("input-topic");
// 构建拓扑
source.to("output-topic");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
无状态操作
KStream<String, String> stream = builder.stream("input-topic");
// 1. 过滤
KStream<String, String> filtered = stream.filter((key, value) -> value.length() > 10);
// 2. filterNot(保留不匹配的)
KStream<String, String> filteredNot = stream.filterNot((key, value) -> value.startsWith("ERROR"));
// 3. 映射键值
KStream<String, Integer> mapped = stream.map((key, value) ->
KeyValue.pair(value.toUpperCase(), value.length()));
// 4. flatMap(一对多)
KStream<String, String> flatMapped = stream.flatMap((key, value) -> {
String[] words = value.toLowerCase().split("\\W+");
List<KeyValue<String, String>> result = new ArrayList<>();
for (String word : words) {
result.add(KeyValue.pair(word, word));
}
return result;
});
// 5. flatMapValues(只展开值)
KStream<String, String> flatMapValues = stream.flatMapValues(value ->
Arrays.asList(value.toLowerCase().split("\\W+")));
// 6. 选择键
KStream<String, String> withNewKey = stream.selectKey((key, value) -> value.split(":")[0]);
// 7. 合并多个流
KStream<String, String> stream1 = builder.stream("topic1");
KStream<String, String> stream2 = builder.stream("topic2");
KStream<String, String> merged = stream1.merge(stream2);
// 8. 分支
KStream<String, String>[] branches = stream.branch(
(key, value) -> value.contains("important"), // 分支 0
(key, value) -> value.contains("warning"), // 分支 1
(key, value) -> true // 分支 2(默认)
);
branches[0].to("important-topic");
branches[1].to("warning-topic");
branches[2].to("other-topic");
有状态操作
// 1. 分组
KGroupedStream<String, String> grouped = stream.groupByKey();
// 按新键分组
KGroupedStream<String, String> groupedByValue = stream.groupBy(
(key, value) -> value.split(":")[0],
Grouped.with(Serdes.String(), Serdes.String())
);
// 2. 聚合 - 计数
KTable<String, Long> counts = grouped.count();
// 持久化状态存储
grouped.count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store")
.withRetention(Duration.ofHours(1))
);
// 3. 聚合 - 求和
KTable<String, Long> sumCounts = stream.groupByKey()
.aggregate(
() -> 0L, // 初始值
(key, value, aggregate) -> aggregate + 1,
Materialized.with(Serdes.String(), Serdes.Long())
);
// 4. 聚合 - 自定义聚合
KTable<String, String> aggregated = stream.groupByKey()
.aggregate(
() -> "", // 初始值
(key, value, aggregate) -> aggregate + value,
Materialized.as("aggregation-store")
);
// 5. 归约
KTable<String, String> reduced = stream.groupByKey()
.reduce(
(value1, value2) -> value1 + "-" + value2,
Materialized.as("reduced-store")
);
// 6. 窗口化操作
// 滚动窗口(大小和滑动间隔相同)
KStream<String, String> windowedStream = builder.stream("input-topic");
KTable<Windowed<String>, Long> rollingCounts = windowedStream
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
.count();
// 跳跃窗口(大小和滑动间隔不同)
KTable<Windowed<String>, Long> hoppingCounts = windowedStream
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)).advanceBy(Duration.ofMinutes(1)))
.count();
// 会话窗口(动态间隔)
KTable<Windowed<String>, Long> sessionCounts = windowedStream
.groupByKey()
.windowedBy(SessionWindows.withGracePeriod(Duration.ofMinutes(5)))
.count();
连接操作
KStream<String, String> leftStream = builder.stream("left-topic");
KStream<String, String> rightStream = builder.stream("right-topic");
// 1. KStream-KStream 连接
KStream<String, String> joined = leftStream.join(
rightStream,
(leftValue, rightValue) -> leftValue + ":" + rightValue,
JoinWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)),
StreamJoined.with(Serdes.String(), Serdes.String(), Serdes.String())
);
// 2. KStream-KTable 连接
KTable<String, String> rightTable = builder.table("right-table-topic");
KStream<String, String> streamTableJoin = leftStream.join(
rightTable,
(leftValue, rightValue) -> leftValue + "-" + rightValue
);
// 3. KStream-KTable 左连接
KStream<String, String> leftJoin = leftStream.leftJoin(
rightTable,
(leftValue, rightValue) -> {
if (rightValue == null) {
return leftValue + "-NULL";
}
return leftValue + "-" + rightValue;
}
);
// 4. KTable-KTable 连接
KTable<String, String> leftTable = builder.table("left-table-topic");
KTable<String, String> tableJoin = leftTable.join(
rightTable,
(leftValue, rightValue) -> leftValue + ":" + rightValue
);
9.3.4 KTable
KTable 是变更日志流的抽象,代表聚合状态。
// 创建 KTable
KTable<String, String> table = builder.table("table-topic");
// 状态存储
KTable<String, Long> counts = stream.groupByKey()
.count(Materialized.as("counts-store"));
// 过滤 KTable
KTable<String, Long> filtered = table.filter((key, value) -> value > 100);
// 映射值
KTable<String, String> mappedValues = table.mapValues(value -> value.toUpperCase());
// 转换为 KStream
KStream<String, Long> countsStream = counts.toStream();
9.3.5 StreamsConfig
Properties props = new Properties();
// 必需配置
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "my-streams-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 序列化配置
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 状态目录
props.put(StreamsConfig.STATE_DIR_CONFIG, "/tmp/kafka-streams");
// 线程数
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 3);
// 缓存大小
props.put(StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG, 10 * 1024 * 1024);
// 提交间隔
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
// Exactly-once 语义
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, "exactly_once_v2");
// 副本数
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);
// 交互式查询
props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "localhost:8080");
9.4 Connect API
9.4.1 概述
Kafka Connect 用于在 Kafka 和其他系统之间传输数据。
9.4.2 连接器配置
{
"name": "my-source-connector",
"config": {
"connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
"tasks.max": "3",
"file": "/tmp/test.txt",
"topic": "connect-test",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.storage.StringConverter"
}
}
9.4.3 REST API
# 列出连接器
curl -X GET http://localhost:8083/connectors
# 创建连接器
curl -X POST http://localhost:8083/connectors \
-H "Content-Type: application/json" \
-d '{"name": "my-connector", "config": {...}}'
# 获取连接器信息
curl -X GET http://localhost:8083/connectors/my-connector
# 获取连接器配置
curl -X GET http://localhost:8083/connectors/my-connector/config
# 获取连接器状态
curl -X GET http://localhost:8083/connectors/my-connector/status
# 暂停连接器
curl -X PUT http://localhost:8083/connectors/my-connector/pause
# 恢复连接器
curl -X PUT http://localhost:8083/connectors/my-connector/resume
# 重启连接器
curl -X POST http://localhost:8083/connectors/my-connector/restart
# 重启任务
curl -X POST http://localhost:8083/connectors/my-connector/tasks/0/restart
# 删除连接器
curl -X DELETE http://localhost:8083/connectors/my-connector
# 获取连接器任务
curl -X GET http://localhost:8083/connectors/my-connector/tasks
9.4.4 SourceConnector 开发
public class MySourceConnector extends SourceConnector {
private String topic;
private String fileName;
@Override
public void start(Map<String, String> props) {
this.topic = props.get("topic");
this.fileName = props.get("file");
}
@Override
public Class<? extends Task> taskClass() {
return MySourceTask.class;
}
@Override
public List<Map<String, String>> taskConfigs(int maxTasks) {
ArrayList<Map<String, String>> configs = new ArrayList<>();
Map<String, String> config = new HashMap<>();
config.put("topic", topic);
config.put("file", fileName);
configs.add(config);
return configs;
}
@Override
public void stop() {
// 清理资源
}
@Override
public ConfigDef config() {
return new ConfigDef()
.define("topic", ConfigDef.Type.STRING, ConfigDef.Importance.HIGH, "目标主题")
.define("file", ConfigDef.Type.STRING, ConfigDef.Importance.HIGH, "源文件");
}
}
9.4.5 SourceTask 开发
public class MySourceTask extends SourceTask {
private String topic;
private BufferedReader reader;
@Override
public void start(Map<String, String> props) {
this.topic = props.get("topic");
String fileName = props.get("file");
try {
reader = new BufferedReader(new FileReader(fileName));
} catch (FileNotFoundException e) {
throw new ConnectException(e);
}
}
@Override
public List<SourceRecord> poll() throws InterruptedException {
try {
ArrayList<SourceRecord> records = new ArrayList<>();
String line;
while ((line = reader.readLine()) != null) {
records.add(new SourceRecord(
Collections.singletonMap("file", fileName),
Collections.singletonMap("offset", lineNum),
topic,
null,
null,
null,
line
));
}
return records.isEmpty() ? null : records;
} catch (IOException e) {
throw new ConnectException(e);
}
}
@Override
public void stop() {
try {
reader.close();
} catch (IOException e) {
// 忽略
}
}
@Override
public String version() {
return "1.0";
}
}
9.5 Admin API
9.5.1 概述
Admin API 用于管理 Kafka 集群,包括创建、删除、修改主题等操作。
9.5.2 创建 AdminClient
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 可选配置
props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, 5000);
props.put(AdminClientConfig.DEFAULT_API_TIMEOUT_MS_CONFIG, 30000);
props.put(AdminClientConfig.RETRIES_CONFIG, 3);
// 使用安全配置
props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
props.put(SaslConfigs.SASL_MECHANISM, "SCRAM-SHA-256");
props.put(SaslConfigs.SASL_JAAS_CONFIG,
"org.apache.kafka.common.security.scram.ScramLoginModule required username=\"user\" password=\"pass\";");
AdminClient admin = AdminClient.create(props);
9.5.3 主题管理
创建主题
// 创建单个主题
String topicName = "my-topic";
int numPartitions = 3;
short replicationFactor = 2;
NewTopic newTopic = new NewTopic(topicName, numPartitions, replicationFactor);
// 设置主题配置
Map<String, String> configs = new HashMap<>();
configs.put("retention.ms", "604800000"); // 7天
configs.put("compression.type", "snappy");
newTopic.configs(configs);
CreateTopicsResult result = admin.createTopics(Collections.singleton(newTopic));
// 等待创建完成
try {
result.all().get();
System.out.println("主题创建成功");
} catch (InterruptedException | ExecutionException e) {
System.err.println("主题创建失败: " + e.getMessage());
}
// 批量创建主题
Map<String, NewTopic> newTopics = new HashMap<>();
newTopics.put("topic1", new NewTopic("topic1", 3, (short) 2));
newTopics.put("topic2", new NewTopic("topic2", 6, (short) 2));
CreateTopicsResult results = admin.createTopics(newTopics);
删除主题
// 删除单个主题
DeleteTopicsResult result = admin.deleteTopics(Collections.singleton("my-topic"));
try {
result.all().get();
System.out.println("主题删除成功");
} catch (InterruptedException | ExecutionException e) {
System.err.println("主题删除失败: " + e.getMessage());
}
// 批量删除
DeleteTopicsResult results = admin.deleteTopics(Arrays.asList("topic1", "topic2"));
列出主题
// 列出所有主题
ListTopicsResult result = admin.listTopics();
Set<String> names = result.names().get();
System.out.println("所有主题: " + names);
// 列出内部主题
ListTopicsOptions options = new ListTopicsOptions();
options.listInternal(true);
Set<String> allNames = admin.listTopics(options).names().get();
描述主题
// 描述单个主题
DescribeTopicsResult describeResult = admin.describeTopics(Collections.singleton("my-topic"));
TopicDescription description = describeResult.allTopics().get().get("my-topic");
System.out.println("主题名称: " + description.name());
System.out.println("分区数: " + description.partitions().size());
for (TopicPartitionInfo partition : description.partitions()) {
System.out.println("分区 " + partition.partition() + ":");
System.out.println(" Leader: " + partition.leader().id());
System.out.println(" Replicas: " + partition.replicas());
System.out.println(" ISR: " + partition.isr());
}
// 批量描述
DescribeTopicsResult results = admin.describeTopics(Arrays.asList("topic1", "topic2"));
修改主题配置
// 修改主题配置
ConfigEntry retentionEntry = new ConfigEntry("retention.ms", "259200000"); // 3天
ConfigEntry compressionEntry = new ConfigEntry("compression.type", "lz4");
AlterConfigOp op1 = new AlterConfigOp(retentionEntry, AlterConfigOp.OpType.SET);
AlterConfigOp op2 = new AlterConfigOp(compressionEntry, AlterConfigOp.OpType.SET);
Map<ConfigResource, Collection<AlterConfigOp>> configs = new HashMap<>();
ConfigResource resource = new ConfigResource(ConfigResource.Type.TOPIC, "my-topic");
configs.put(resource, Arrays.asList(op1, op2));
AlterConfigsResult alterResult = admin.incrementalAlterConfigs(configs);
alterResult.all().get();
// 删除配置(恢复默认值)
AlterConfigOp deleteOp = new AlterConfigOp(
new ConfigEntry("retention.ms", ""),
AlterConfigOp.OpType.DELETE
);
Map<ConfigResource, Collection<AlterConfigOp>> deleteConfigs = new HashMap<>();
deleteConfigs.put(resource, Collections.singleton(deleteOp));
admin.incrementalAlterConfigs(deleteConfigs).all().get();
修改分区数
// 增加分区数(只能增加,不能减少)
Map<String, NewPartitions> newPartitions = new HashMap<>();
newPartitions.put("my-topic", NewPartitions.increaseTo(10));
CreatePartitionsResult result = admin.createPartitions(newPartitions);
result.all().get();
9.5.4 配置管理
描述配置
// 描述主题配置
ConfigResource topicResource = new ConfigResource(ConfigResource.Type.TOPIC, "my-topic");
DescribeConfigsResult describeResult = admin.describeConfigs(Collections.singleton(topicResource));
Config config = describeResult.all().get().get(topicResource);
for (ConfigEntry entry : config.entries()) {
System.out.println(entry.name() + " = " + entry.value() +
" (来源: " + entry.source() + ")");
}
描述 Broker 配置
// 描述 Broker 配置
ConfigResource brokerResource = new ConfigResource(ConfigResource.Type.BROKER, "0");
DescribeConfigsResult describeResult = admin.describeConfigs(Collections.singleton(brokerResource));
Config config = describeResult.all().get().get(brokerResource);
9.5.5 消费者组管理
列出消费者组
// 列出所有消费者组
ListConsumerGroupsResult result = admin.listConsumerGroups();
Collection<ConsumerGroupListing> groups = result.all().get();
for (ConsumerGroupListing group : groups) {
System.out.println("组ID: " + group.groupId());
System.out.println("是否简单: " + group.isSimpleConsumerGroup());
}
描述消费者组
// 描述消费者组
DescribeConsumerGroupsResult describeResult = admin.describeConsumerGroups(
Collections.singleton("my-group")
);
ConsumerGroupDescription groupDescription = describeResult.all().get().get("my-group");
System.out.println("状态: " + groupDescription.state());
System.out.println("成员数: " + groupDescription.members().size());
System.out.println("分区分配策略: " + groupDescription.partitionAssignor());
for (MemberDescription member : groupDescription.members()) {
System.out.println("成员: " + member.consumerId());
System.out.println(" 客户端ID: " + member.clientId());
System.out.println(" 主机: " + member.host());
System.out.println(" 分配的分区: " + member.assignment().topicPartitions());
}
获取消费者组偏移量
// 获取消费者组的偏移量
ListConsumerGroupOffsetsResult offsetsResult = admin.listConsumerGroupOffsets("my-group");
Map<TopicPartition, OffsetAndMetadata> offsets = offsetsResult.partitionsToOffsetAndMetadata().get();
for (Map.Entry<TopicPartition, OffsetAndMetadata> entry : offsets.entrySet()) {
System.out.println(entry.getKey() + ": " + entry.getValue().offset());
}
删除消费者组
// 删除消费者组(必须是空的)
DeleteConsumerGroupsResult result = admin.deleteConsumerGroups(
Collections.singleton("my-group")
);
try {
result.all().get();
System.out.println("消费者组删除成功");
} catch (ExecutionException e) {
System.err.println("删除失败: " + e.getMessage());
}
重置偏移量
// 重置消费者组偏移量到最早
Map<TopicPartition, OffsetAndMetadata> resetOffsets = new HashMap<>();
TopicPartition partition = new TopicPartition("my-topic", 0);
resetOffsets.put(partition, new OffsetAndMetadata(0)); // 或者使用 OffsetSpec.earliest()
AlterConsumerGroupOffsetsResult result = admin.alterConsumerGroupOffsets(
"my-group",
resetOffsets
);
result.all().get();
9.5.6 集群管理
获取集群信息
// 获取集群信息
DescribeClusterResult clusterResult = admin.describeCluster();
// 获取集群ID
String clusterId = clusterResult.clusterId().get();
// 获取 Controller
Node controller = clusterResult.controller().get();
// 获取 Broker 列表
Collection<Node> brokers = clusterResult.nodes().get();
for (Node broker : brokers) {
System.out.println("Broker " + broker.id() + ": " + broker.host() + ":" + broker.port());
}
// 获取授权操作
AclOperationOperations operations = clusterResult.authorizedOperations().get();
获取分区信息
// 查询主题的可用分区
Map<String, TopicDescription> topicDescriptions = admin.describeTopics(
Collections.singleton("my-topic")
).allTopicNames().get();
for (TopicPartitionInfo partition : topicDescriptions.get("my-topic").partitions()) {
TopicPartition tp = new TopicPartition("my-topic", partition.partition());
// 处理分区信息
}
9.5.7 ACL 管理
创建 ACL
// 创建主题 ACL
AclBinding aclBinding = new AclBinding(
new Resource(ResourceType.TOPIC, "my-topic"),
new AccessControlEntry("User:alice", "*", AclOperation.READ, AclPermissionType.ALLOW)
);
CreateAclsResult result = admin.createAcls(Collections.singleton(aclBinding));
result.all().get();
列出 ACL
// 列出主题的所有 ACL
ResourcePatternFilter filter = new ResourcePatternFilter(
ResourceType.TOPIC,
"my-topic",
PatternType.MATCH
);
DescribeAclsResult result = admin.describeAcls(new AclBindingFilter(filter, AccessControlEntryFilter.ANY));
Collection<AclBinding> acls = result.values().get();
for (AclBinding acl : acls) {
System.out.println(acl);
}
删除 ACL
// 删除 ACL
AclBindingFilter filter = new AclBindingFilter(
new ResourcePatternFilter(ResourceType.TOPIC, "my-topic", PatternType.EXACT),
new AccessControlEntryFilter("User:alice", "*", AclOperation.READ, AclPermissionType.ALLOW)
);
DeleteAclsResult result = admin.deleteAcls(Collections.singleton(filter));
result.all().get();
9.6 交互式查询
9.6.1 概述
交互式查询允许从 Kafka Streams 应用外部查询应用的状态。
9.6.2 查询本地状态存储
查询键值存储
// 获取键值存储
ReadOnlyKeyValueStore<String, Long> keyValueStore = streams.store(
"CountsKeyValueStore",
QueryableStoreTypes.keyValueStore()
);
// 通过键获取值
Long count = keyValueStore.get("hello");
// 获取范围
KeyValueIterator<String, Long> range = keyValueStore.range("a", "z");
while (range.hasNext()) {
KeyValue<String, Long> next = range.next();
System.out.println(next.key + ": " + next.value);
}
// 获取所有记录
KeyValueIterator<String, Long> all = keyValueStore.all();
while (all.hasNext()) {
KeyValue<String, Long> next = all.next();
System.out.println(next.key + ": " + next.value);
}
// 获取近似条目数
long approximateNumEntries = keyValueStore.approximateNumEntries();
查询窗口存储
// 获取窗口存储
ReadOnlyWindowStore<String, Long> windowStore = streams.store(
"CountsWindowStore",
QueryableStoreTypes.windowStore()
);
// 获取特定键的时间窗口
Instant timeFrom = Instant.ofEpochMilli(0); // 开始时间
Instant timeTo = Instant.now(); // 结束时间
WindowStoreIterator<Long> iterator = windowStore.fetch("world", timeFrom, timeTo);
while (iterator.hasNext()) {
KeyValue<Long, Long> next = iterator.next();
long windowTimestamp = next.key;
long count = next.value;
System.out.println("时间 " + windowTimestamp + ": " + count);
}
// 获取所有键的最新窗口
KeyValueIterator<Windowed<String>, Long> all = windowStore.all();
while (all.hasNext()) {
KeyValue<Windowed<String>, Long> next = all.next();
System.out.println(next.key + ": " + next.value);
}
9.6.3 查询远程状态存储
配置 RPC 端点
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_SERVER_CONFIG, "host1:4460");
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> textLines = builder.stream("word-count-input");
KGroupedStream<String, String> groupedByWord = textLines
.flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
.groupBy((key, word) -> word);
// 创建状态存储
groupedByWord.count(Materialized.as("word-count"));
KafkaStreams streams = new KafkaStreams(builder, props);
streams.start();
// 启动 RPC 服务
// MyRPCService rpcService = ...;
// rpcService.listenAt("host1:4460");
发现远程应用实例
KafkaStreams streams = ...;
// 查找所有包含特定状态存储的应用实例
Collection<StreamsMetadata> metadataList = streams.allMetadataForStore("word-count");
for (StreamsMetadata metadata : metadataList) {
System.out.println("主机: " + metadata.host() + ":" + metadata.port());
System.out.println("状态存储: " + metadata.stateStoreNames());
System.out.println("分区: " + metadata.topicPartitions());
}
// 查找包含特定键的应用实例
StreamsMetadata metadata = streams.metadataForKey(
"word-count",
"alice",
Serdes.String().serializer()
);
if (metadata != null) {
String host = metadata.host();
int port = metadata.port();
// 通过 RPC 查询远程实例
}
// 获取所有应用实例的元数据
Collection<StreamsMetadata> allMetadata = streams.allMetadata();
RPC 查询示例
// 查询特定键的方法1:使用 metadataForKey
StreamsMetadata metadata = streams.metadataForKey(
"word-count",
"alice",
Serdes.String().serializer()
);
if (metadata != null) {
String url = "http://" + metadata.host() + ":" + metadata.port() + "/word-count/alice";
Long result = httpClient.getLong(url);
}
// 查询特定键的方法2:遍历所有实例
Optional<Long> result = streams.allMetadataForStore("word-count")
.stream()
.map(streamsMetadata -> {
String url = "http://" + streamsMetadata.host() + ":" +
streamsMetadata.port() + "/word-count/alice";
return httpClient.getLong(url);
})
.filter(s -> s != null)
.findFirst();
9.7 应用重置工具
9.7.1 概述
应用重置工具允许你重置 Kafka Streams 应用,强制它从头开始重新处理数据。
9.7.2 工作原理
重置工具对不同类型的主题采取不同操作:
| 主题类型 | 操作 |
|---|---|
| 输入主题 | 根据参数重置偏移量到最早、最晚等位置 |
| 输出主题 | 不做任何操作(不删除数据) |
| 中间主题 | 跳到末尾(跳过所有现有消息) |
| 内部主题 | 删除并重新创建 |
9.7.3 使用前准备
# 1. 停止所有应用实例
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-application-id
# 2. 验证消费者组是否不再活动
# 3. (可选)删除并重新创建中间主题以释放磁盘空间
9.7.4 重置工具参数
$ bin/kafka-streams-application-reset
Option (* = required) Description
--------------------- -----------
* --application-id <String: id> Kafka Streams 应用 ID (application.id)
--bootstrap-server <String: urls> Kafka 服务器地址
--by-duration <String: duration> 按持续时间重置偏移量,格式: 'PnDTnHnMnS'
--dry-run 预览将要执行的操作
--from-file <String: file> 从 CSV 文件读取偏移量值
--input-topics <String: list> 输入主题列表(逗号分隔)
--intermediate-topics <String: list> 中间主题列表(逗号分隔)
--internal-topics <String: list> 内部主题列表(逗号分隔)
--shift-by <Long: n> 将当前偏移量移动 n 位
--to-datetime <String> 重置到指定日期时间,格式: 'YYYY-MM-DDTHH:mm:SS.sss'
--to-earliest 重置到最早偏移量
--to-latest 重置到最晚偏移量
--to-offset <Long> 重置到特定偏移量
--force 强制移除消费者组成员
9.7.5 重置场景
场景1:从头重新处理所有数据
# 第一步:运行应用重置工具
bin/kafka-streams-application-reset.sh \
--application-id my-app \
--bootstrap-server localhost:9092 \
--input-topics input-topic-1,input-topic-2 \
--intermediate-topics intermediate-topic-1,intermediate-topic-2 \
--to-earliest
# 第二步:删除本地状态目录
rm -rf /tmp/kafka-streams/my-app
# 第三步:重启应用
场景2:跳过已有数据,只处理新数据
bin/kafka-streams-application-reset.sh \
--application-id my-app \
--bootstrap-server localhost:9092 \
--input-topics input-topic-1 \
--to-latest
rm -rf /tmp/kafka-streams/my-app
场景3:重置到特定时间点
# 重置到24小时前
bin/kafka-streams-application-reset.sh \
--application-id my-app \
--bootstrap-server localhost:9092 \
--input-topics input-topic-1 \
--by-duration P1D
rm -rf /tmp/kafka-streams/my-app
场景4:使用 API 清理
KafkaStreams streams = ...;
// 在关闭应用后调用
// 清理本地状态目录
streams.cleanUp();
// 或手动删除
// Files.delete(Paths.get("/tmp/kafka-streams/my-app"));
9.7.6 注意事项
- 停止所有实例:所有应用实例必须停止
- 检查参数:错误参数可能影响其他应用
- 中间主题:建议手动删除并重建以释放空间
- 长期会话超时:如果配置了长会话超时,使用
--force选项 - 不可逆操作:建议先使用
--dry-run预览
附录:API 类图
Producer API 类图
┌─────────────────┐
│ KafkaProducer │
├─────────────────┤
│ + send() │
│ + flush() │
│ + close() │
│ + initTransactions() │
│ + commitTransaction() │
└────────┬────────┘
│
│ creates
▼
┌─────────────────┐
│ ProducerRecord │
├─────────────────┤
│ - topic │
│ - partition │
│ - key │
│ - value │
│ - timestamp │
│ - headers │
└─────────────────┘
Consumer API 类图
┌─────────────────┐
│ KafkaConsumer │
├─────────────────┤
│ + subscribe() │
│ + assign() │
│ + poll() │
│ + commitSync() │
│ + commitAsync() │
│ + seek() │
│ + position() │
│ + committed() │
└────────┬────────┘
│
│ returns
▼
┌─────────────────┐
│ ConsumerRecords │
├─────────────────┤
│ + records() │
│ + partitions() │
│ + count() │
│ + isEmpty() │
└────────┬────────┘
│
│ contains
▼
┌─────────────────┐
│ ConsumerRecord │
├─────────────────┤
│ - topic │
│ - partition │
│ - offset │
│ - key │
│ - value │
│ - timestamp │
│ - headers │
└─────────────────┘
Streams API 类图
┌─────────────────┐
│ StreamsBuilder │
├─────────────────┤
│ + stream() │ → KStream
│ + table() │ → KTable
│ + globalTable() │ → GlobalKTable
│ + build() │ → Topology
└─────────────────┘
┌─────────────────┐
│ KStream │
├─────────────────┤
│ + filter() │
│ + map() │
│ + flatMap() │
│ + groupByKey() │
│ + join() │
│ + to() │
└─────────────────┘
┌─────────────────┐
│ KTable │
├─────────────────┤
│ + filter() │
│ + mapValues() │
│ + join() │
│ + toStream() │
└─────────────────┘