1. Kafka 自动发送消息 Demo 概述
在分布式系统架构中,消息队列扮演着至关重要的角色。Kafka 作为一款高性能、高吞吐量的分布式消息系统,已经成为现代互联网企业的基础设施标配。这个 Demo 将展示如何用 Java 语言实现 Kafka 消息的自动发送功能,涵盖从环境配置到代码实现的完整流程。
对于刚接触 Kafka 的开发者来说,第一个需要攻克的难关就是如何正确地配置和发送消息。很多新手在初次尝试时容易陷入各种配置陷阱,比如连接不上 broker、消息发送失败却无报错等问题。本文将基于实战经验,带你避开这些常见坑点。
2. 环境准备与配置
2.1 Kafka 服务端安装
首先需要搭建 Kafka 服务端环境。推荐使用最新稳定版本(当前为 3.5.0),可以从 Apache 官网下载二进制包。解压后目录结构包含:
- bin/: 各种可执行脚本
- config/: 配置文件目录
- libs/: 依赖库
启动 Kafka 前需要先启动 Zookeeper(单机开发环境可以使用 Kafka 内置的 Zookeeper):
# 启动 Zookeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 启动 Kafka broker bin/kafka-server-start.sh config/server.properties注意:生产环境建议使用外置 Zookeeper 集群,并配置多个 broker 节点实现高可用。
2.2 Java 项目依赖配置
在 Maven 项目中添加 Kafka 客户端依赖:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.5.0</version> </dependency>如果是 Gradle 项目:
implementation 'org.apache.kafka:kafka-clients:3.5.0'3. 生产者配置详解
3.1 核心配置参数
创建 KafkaProducer 时需要配置一些必要参数:
Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); // broker地址 props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("acks", "all"); // 消息确认机制 props.put("retries", 3); // 重试次数 props.put("linger.ms", 5); // 发送延迟关键参数说明:
| 参数 | 说明 | 推荐值 |
|---|---|---|
| bootstrap.servers | broker地址列表 | 生产环境建议配置多个 |
| acks | 消息确认机制 | all(最安全) |
| retries | 发送失败重试次数 | 3-5 |
| batch.size | 批量发送大小 | 16384-65536 |
| linger.ms | 发送等待时间 | 5-100 |
3.2 序列化器选择
Kafka 消息的 key 和 value 都需要指定序列化器。除了内置的 StringSerializer,还可以使用:
- ByteArraySerializer
- IntegerSerializer
- JSON 序列化(如 Jackson)
- Avro 序列化
对于复杂对象,推荐使用 JSON 或 Avro 格式:
props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer"); // 使用 Jackson 将对象转为 JSON bytes ObjectMapper mapper = new ObjectMapper(); byte[] jsonBytes = mapper.writeValueAsBytes(myObject);4. 消息发送实战
4.1 基础发送模式
创建生产者并发送消息的基本流程:
KafkaProducer<String, String> producer = new KafkaProducer<>(props); try { for(int i = 0; i < 100; i++) { ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "key-" + i, "value-" + i); // 同步发送 RecordMetadata metadata = producer.send(record).get(); System.out.printf("Sent record(key=%s value=%s) to partition=%d offset=%d%n", record.key(), record.value(), metadata.partition(), metadata.offset()); } } finally { producer.close(); }4.2 异步发送与回调
为提高吞吐量,通常使用异步发送方式:
producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception e) { if(e != null) { log.error("Send failed for record {}", record, e); } else { log.debug("Sent to {}-{}@{}", metadata.topic(), metadata.partition(), metadata.offset()); } } });4.3 消息分区策略
Kafka 通过分区实现并行处理。指定分区的方式有:
- 显式指定分区号
- 通过 key 的 hash 计算分区
- 自定义分区器
// 1. 直接指定分区 new ProducerRecord<>("topic", 0, "key", "value"); // 2. 使用 key 的 hash(默认) new ProducerRecord<>("topic", "key", "value"); // 3. 自定义分区器 props.put("partitioner.class", "com.my.CustomPartitioner");5. 高级特性与优化
5.1 事务消息
Kafka 支持跨分区的事务操作:
props.put("enable.idempotence", "true"); props.put("transactional.id", "my-transactional-id"); producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord<>("topic1", "key", "value")); producer.send(new ProducerRecord<>("topic2", "key", "value")); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }5.2 消息压缩
为减少网络传输量,可以启用压缩:
props.put("compression.type", "snappy"); // 或 gzip, lz4压缩算法对比:
| 算法 | 压缩率 | 速度 | CPU消耗 |
|---|---|---|---|
| gzip | 高 | 慢 | 高 |
| snappy | 中 | 快 | 中 |
| lz4 | 低 | 最快 | 低 |
5.3 性能调优
提升发送性能的关键参数:
props.put("buffer.memory", 33554432); // 缓冲区大小 props.put("max.block.ms", 60000); // 阻塞超时 props.put("request.timeout.ms", 30000); // 请求超时6. 问题排查与监控
6.1 常见问题排查
连接失败:
- 检查防火墙设置
- 确认 broker 地址正确
- 检查网络连通性
消息发送失败:
- 检查 topic 是否存在
- 查看 broker 日志
- 调整重试策略
性能低下:
- 增加批量大小
- 调整 linger.ms
- 启用压缩
6.2 监控指标
关键监控指标包括:
- 请求速率
- 请求延迟
- 批量大小
- 错误率
可以使用 JMX 或 Prometheus 收集这些指标:
props.put("metric.reporters", "com.my.MetricsReporter"); props.put("metrics.num.samples", "2"); props.put("metrics.sample.window.ms", "30000");7. 完整示例代码
下面是一个完整的自动发送消息示例:
public class KafkaAutoProducer { private static final Logger log = LoggerFactory.getLogger(KafkaAutoProducer.class); private volatile boolean running = true; public void start(String topic, long interval) { 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"); // 性能优化 props.put("linger.ms", "5"); props.put("batch.size", "16384"); props.put("compression.type", "snappy"); KafkaProducer<String, String> producer = new KafkaProducer<>(props); Runtime.getRuntime().addShutdownHook(new Thread(() -> { running = false; producer.close(); })); int count = 0; while(running) { try { String key = "key-" + (count % 10); String value = "value-" + System.currentTimeMillis(); ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value); producer.send(record, (metadata, e) -> { if(e != null) { log.error("Send failed", e); } else { log.debug("Sent to {}-{}@{}", metadata.topic(), metadata.partition(), metadata.offset()); } }); count++; Thread.sleep(interval); } catch (Exception e) { log.error("Error in producer", e); } } } }8. 生产环境建议
资源隔离:
- 为 Kafka 分配专用服务器
- 生产者和消费者使用独立的网络带宽
容错处理:
- 实现消息重试机制
- 添加死信队列处理
- 监控关键指标
安全配置:
- 启用 SSL 加密
- 配置 SASL 认证
- 设置 ACL 权限控制
// 安全配置示例 props.put("security.protocol", "SASL_SSL"); props.put("sasl.mechanism", "PLAIN"); props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required " + "username=\"user\" password=\"pwd\";");9. 性能测试与优化
9.1 基准测试
使用 kafka-producer-perf-test 工具进行测试:
bin/kafka-producer-perf-test.sh \ --topic test \ --num-records 1000000 \ --record-size 1000 \ --throughput -1 \ --producer-props \ bootstrap.servers=localhost:9092 \ batch.size=16384 \ linger.ms=09.2 优化方向
根据测试结果可能的优化点:
- 增加批量大小(batch.size)
- 调整等待时间(linger.ms)
- 启用压缩(compression.type)
- 增加生产者实例数
- 优化网络配置
10. 与其他系统集成
10.1 Spring Kafka 集成
Spring Boot 提供了便捷的 Kafka 集成:
@Configuration public class KafkaConfig { @Bean public ProducerFactory<String, String> producerFactory() { Map<String, Object> config = new HashMap<>(); config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return new DefaultKafkaProducerFactory<>(config); } @Bean public KafkaTemplate<String, String> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } } @Service public class MessageService { @Autowired private KafkaTemplate<String, String> kafkaTemplate; public void send(String topic, String message) { kafkaTemplate.send(topic, message); } }10.2 与流处理系统集成
Kafka 消息可以被 Flink、Spark Streaming 等系统消费:
// Flink 消费 Kafka 示例 FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>( "input-topic", new SimpleStringSchema(), properties); DataStream<String> stream = env.addSource(consumer);