跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 10 / 10 查看系列简介 →

Kafka 3.9 教程(九):API 参考

文章目录
  1. 目录
  2. 9.1 Producer API
  3. 9.1.1 概述
  4. 9.1.2 Maven 依赖
  5. 9.1.3 ProducerRecord
  6. 9.1.4 KafkaProducer
  7. 9.1.5 RecordMetadata
  8. 9.1.6 Producer 拦截器
  9. 9.2 Consumer API
  10. 9.2.1 概述
  11. 9.2.2 KafkaConsumer
  12. 9.2.3 Consumer 拦截器
  13. 9.3 Streams API
  14. 9.3.1 概述
  15. 9.3.2 Maven 依赖
  16. 9.3.3 KStream
  17. 9.3.4 KTable
  18. 9.3.5 StreamsConfig
  19. 9.4 Connect API
  20. 9.4.1 概述
  21. 9.4.2 连接器配置
  22. 9.4.3 REST API
  23. 9.4.4 SourceConnector 开发
  24. 9.4.5 SourceTask 开发
  25. 9.5 Admin API
  26. 9.5.1 概述
  27. 9.5.2 创建 AdminClient
  28. 9.5.3 主题管理
  29. 9.5.4 配置管理
  30. 9.5.5 消费者组管理
  31. 9.5.6 集群管理
  32. 9.5.7 ACL 管理
  33. 9.6 交互式查询
  34. 9.6.1 概述
  35. 9.6.2 查询本地状态存储
  36. 9.6.3 查询远程状态存储
  37. 9.7 应用重置工具
  38. 9.7.1 概述
  39. 9.7.2 工作原理
  40. 9.7.3 使用前准备
  41. 9.7.4 重置工具参数
  42. 9.7.5 重置场景
  43. 9.7.6 注意事项
  44. 附录:API 类图
  45. Producer API 类图
  46. Consumer API 类图
  47. Streams API 类图
  48. 下一步学习

📌 本文是「Kafka 3.9 教程」系列第 九 篇 · 系列目录

本文档基于 Apache Kafka 3.9 官方文档整理

目录


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 注意事项

  1. 停止所有实例:所有应用实例必须停止
  2. 检查参数:错误参数可能影响其他应用
  3. 中间主题:建议手动删除并重建以释放空间
  4. 长期会话超时:如果配置了长会话超时,使用 --force 选项
  5. 不可逆操作:建议先使用 --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()    │
└─────────────────┘

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(九):API 参考》