Skip to content

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 构建实时数据管道
  • 事件溯源:记录系统中所有状态变更事件
  • 指标监控:收集服务和基础设施的监控数据