Java Kafka生产者消费者案例

wen java案例 3

本文目录导读:

Java Kafka生产者消费者案例

  1. 目录导读
  2. Kafka核心概念速览
  3. 环境准备与依赖配置
  4. Java生产者案例
  5. Java消费者案例
  6. 性能调优与常见坑
  7. 实战问答(FAQ)
  8. 总结与最佳实践

Java Kafka生产者消费者实战:从入门到性能调优的完整指南

目录导读

  1. Kafka核心概念速览:为什么需要消息队列?Kafka的架构角色与术语
  2. 环境准备与依赖配置:Maven依赖、Kafka服务端启动参数
  3. Java生产者案例:同步/异步发送、分区策略、重试机制、回调处理
  4. Java消费者案例:订阅模式、位移提交(自动/手动)、消费者组与再均衡
  5. 性能调优与常见坑:批量大小、压缩、拦截器、序列化器陷阱
  6. 实战问答(FAQ):解决开发者高频疑问
  7. 总结与最佳实践:生产环境推荐配置清单

Kafka核心概念速览

Apache Kafka是一个分布式流处理平台,其核心能力是高吞吐、低延迟的消息发布与订阅,在开始编写Java代码之前,理解以下核心术语至关重要:

  • Broker:Kafka集群中的一台服务器节点。
  • Topic:消息的逻辑分类,类似数据库中的表。
  • Partition:Topic的物理分片,每个分区内部有序,分区之间可以并行读写。
  • Offset:消息在分区内的唯一序号,消费者依靠它记录消费位置。
  • Consumer Group:同一组内的消费者共享一个Topic的消息,实现负载均衡;不同组之间互不影响。

架构亮点:Kafka通过“顺序写磁盘”和“零拷贝”技术,实现了单机百万级消息吞吐,这也是它相比RabbitMQ、RocketMQ在日志聚合、用户行为追踪场景中更占优势的原因。


环境准备与依赖配置

1 启动Kafka服务端(本地开发)

# 启动Zookeeper(Kafka 2.8+支持KRaft模式无需ZK,此处用经典模式)
bin/zookeeper-server-start.sh config/zookeeper.properties
# 启动Broker
bin/kafka-server-start.sh config/server.properties
# 创建测试Topic:3个分区,2个副本
bin/kafka-topics.sh --create --topic order-topic --partitions 3 --replication-factor 2 --bootstrap-server localhost:9092

2 Maven依赖(使用最新稳定版本3.x)

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-clients</artifactId>
    <version>3.6.0</version>
</dependency>

所有API均在org.apache.kafka.clients包下,无需额外引入第三方库。


Java生产者案例

1 基本配置与同步发送

import org.apache.kafka.clients.producer.*;
import java.util.Properties;
public class SimpleProducer {
    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        // 指定Broker地址,多个用逗号分隔
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        // 指定序列化器,必须与消息类型一致
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer");
        // 等待所有副本确认,提高可靠性(0=不确认,1=Leader确认,all=-1=全部ISR确认)
        props.put(ProducerConfig.ACKS_CONFIG, "all");
        Producer<String, String> producer = new KafkaProducer<>(props);
        // 同步发送:每次send后立即get()等待结果
        for (int i = 0; i < 10; i++) {
            ProducerRecord<String, String> record = new ProducerRecord<>("order-topic", "key-" + i, "value-" + i);
            RecordMetadata metadata = producer.send(record).get();  // 阻塞直到返回
            System.out.printf("发送成功: 分区=%d, 偏移量=%d%n", metadata.partition(), metadata.offset());
        }
        producer.close();  // 必须关闭以释放资源
    }
}

关键点acks=all代表最强一致性,但吞吐量会下降,若允许丢消息,可设为1

2 异步发送与回调(生产推荐)

// 异步发送配合Callback,避免主线程阻塞
producer.send(record, (metadata, exception) -> {
    if (exception == null) {
        System.out.printf("异步发送成功: 主题=%s, 分区=%d, 偏移量=%d%n",
                metadata.topic(), metadata.partition(), metadata.offset());
    } else {
        exception.printStackTrace();  // 实际应接入日志或重试队列
    }
});
// 注意:异步时producer.close()会等待所有消息发送完成,但若回调异常需自行处理。

性能对比:异步+回调比同步快至少10倍(在500条/秒压力测试下),生产环境务必使用异步。

3 自定义分区器与拦截器

// 分区器:按业务字段(如用户ID)取模分区,保证同一用户消息有序
props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, UserIdPartitioner.class.getName());
public class UserIdPartitioner implements Partitioner {
    @Override
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
        Integer partitions = cluster.partitionCountForTopic(topic);
        return Math.abs(key.hashCode()) % partitions;
    }
}
// 拦截器:统计发送成功率、修改消息头等

Java消费者案例

1 自动提交位移(简单但不推荐生产)

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-consumer-group");  // 同一个GroupId才可负载均衡
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");  // 默认true
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");  // 每秒自动提交
Consumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Collections.singletonList("order-topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
    for (ConsumerRecord<String, String> record : records) {
        System.out.printf("收到消息: 分区=%d, offset=%d, key=%s, value=%s%n",
                record.partition(), record.offset(), record.key(), record.value());
    }
    // 不需要手动调用commitSync()
}

风险:自动提交存在“重复消费”或“消息丢失”窗口,例如程序在提交前崩溃,重启后会重复消费本次poll的数据;若在poll()之前就提交了offset,但业务处理失败,则消息丢失。

2 手动提交位移(生产标准)

props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");  // 关闭自动提交
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> record : records) {
        process(record);  // 业务处理:写DB、调用API等
    }
    // 方式1:同步提交(等待成功,但会阻塞)
    consumer.commitSync();  
    // 方式2:异步提交(不阻塞,但重试可能乱序,推荐配合回调)
    // consumer.commitAsync((offsets, exception) -> { if (exception != null) log.error(...); });
}

最佳实践:使用commitSync处理最后一批消息,commitAsync处理中间批次,避免阻塞。

3 消费者组与再均衡监听

consumer.subscribe(Collections.singletonList("order-topic"), new ConsumerRebalanceListener() {
    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        // 再均衡前保存当前offset到DB
        saveOffsetsToDB(consumer.position(partition));
    }
    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // 从DB恢复offset后seek到指定位置
        partitions.forEach(p -> consumer.seek(p, getOffsetFromDB(p)));
    }
});

注意:再均衡监听是防止“重复消费”和“跳过消息”的核心,当消费者数量增减或Topic分区数变更时触发。


性能调优与常见坑

1 生产者优化三宝

参数名 默认值 建议值 说明
batch.size 16KB 32KB-64KB 批量发送字节数,越大吞吐越高,但延迟增加
linger.ms 0 5-20ms 等待更多消息组成批次的时间,配合batch.size使用
compression.type none lz4或zstd 开启压缩可减少网络IO,CPU开销可忽略

2 消费者极端场景调优

  • 单线程处理慢:通过max.poll.interval.ms(默认300秒)控制poll()间隔,若业务处理超时,消费者会被踢出组并触发再均衡。
  • 大量堆积消息:增加fetch.min.bytes(默认1字节)和fetch.max.wait.ms(默认500ms),让服务端攒满数据再返回,减少网络交互次数。
  • 顺序消费:只用一个分区、一个消费者,或自定义分区器保证相同key进同一分区。

3 常见坑警示

  • 序列化器不一致:生产者用JSON,消费者用String,直接反序列化会报SerializationException
  • 消费组ID重复:不同业务必须用不同group.id,否则A业务会消费B业务的消息。
  • 未关闭资源:不调用producer.close()consumer.wakeup()(配合close()),会导致连接泄漏和后台线程卡死。

实战问答(FAQ)

Q1:Kafka只保证分区内有序,如何保证全局有序? A:全局有序必须让Topic只有一个Partition,但这牺牲了并行度,折中方案是按业务主键分区(如订单ID),保证同一订单的消息有序即可。

Q2:消费者“重复消费”如何规避? A:使用手动提交+业务幂等(例如写入数据库时用唯一索引),在onPartitionsRevoked回调中保存position,恢复时seek到已处理位置,可最大限度避免重复。

Q3:生产环境Broker集群崩溃,消息会丢失吗? A:若acks=allmin.insync.replicas=2(至少2个副本同步),则不会丢失,但若只设置了acks=1,Leader单独挂掉时会丢消息,还需要注意unclean.leader.election.enable=false禁止非ISR副本竞选。

Q4:为什么消费者组内的消费者数量超过分区数时,多余的消费者会闲置? A:因为每一个分区在同一时间只能被一个消费者消费,假设5个分区,7个消费者,那么只有5个消费者活跃,2个消费者完全空闲且不消费任何消息,这是Kafka设计使然,不能动态增加分区数来生效(分区数只能增加)。

Q5:KafkaProducer是线程安全吗? A:是线程安全的,可以在多个线程中共享同一个Producer实例,但要注意send()方法是异步的,回调线程(Sender线程)与业务线程是隔离的,不要阻塞回调。


总结与最佳实践

场景 配置推荐
日志/埋点数据(允许少量丢失) acks=1、压缩lz4、批量32KB、异步
订单/支付(要求不丢失) acks=allmin.insync.replicas=2、手动提交+幂等
实时推荐/风控(低延迟) linger.ms=0、单分区单消费者、关压缩

核心心法

  • 生产者:异步大于同步,但必须处理好回调异常;合理配置batch和linger才能发挥Kafka性能。
  • 消费者:手动提交+幂等处理是避免数据错误的不二法门;学会使用再均衡监听器来动态恢复现场。
  • 监控:通过JMX监控request-latency-avgrecords-lag-max等指标,建议接入Prometheus+Grafana。

无论您是初学者还是资深开发者,建议在测试环境模拟消费者宕机、Broker重启等故障场景,亲自验证位移提交机制——这才是掌握Kafka精髓的必经之路。


(全文结束)

抱歉,评论功能暂时关闭!