Kafka
简介
Apache Kafka 是分布式流处理平台,由 LinkedIn 开发并贡献给 Apache 基金会。以高吞吐量、持久化存储和强大的流处理能力著称。
核心概念
- Producer:消息生产者
- Consumer:消息消费者
- Topic:主题,消息的分类
- Partition:分区,Topic 的并行处理单元
- Broker:Kafka 服务器节点
- Consumer Group:消费者组,实现负载均衡和容错
- Offset:偏移量,记录消费者消费的位置
架构特点
Producer1 ───┐
├──> Broker1 (TopicA-P0) ───> Consumer Group A
Producer2 ───┘ ├──> Consumer Group B
Broker2 (TopicA-P1) ───┘高性能设计
- 顺序写入:消息追加到文件末尾,磁盘顺序写性能极高
- 零拷贝:利用 sendfile 系统调用,减少数据在内核和用户态之间拷贝
- 批量处理:消息批量发送、批量压缩、批量刷盘
- 分区并行:分区是 Kafka 并行处理的基本单元
生产与消费
java
// 生产者配置
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("my-topic", "key", "value"));
// 消费者配置
props.put("group.id", "my-group");
props.put("enable.auto.commit", "true");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset = %d, key = %s, value = %s%n",
record.offset(), record.key(), record.value());
}
}适用场景
- 日志聚合:收集和传输分布式系统日志
- 流处理:基于 Kafka Streams 或 Flink 构建实时数据管道
- 事件溯源:记录系统中所有状态变更事件
- 指标监控:收集服务和基础设施的监控数据