本文目录导读:

我来详细介绍Java生产者发送消息的几种常见案例,以Kafka为例:
基础配置和生产消息
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
public class KafkaProducerDemo {
public static void main(String[] args) {
// 1. 配置生产者属性
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());
// 2. 创建生产者对象
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// 3. 发送消息
try {
// 创建消息记录
ProducerRecord<String, String> record = new ProducerRecord<>(
"my-topic", // topic名称
"key-1", // key
"Hello Kafka!" // value
);
// 发送消息(异步)
producer.send(record);
} catch (Exception e) {
e.printStackTrace();
} finally {
// 4. 关闭生产者
producer.close();
}
}
}
发送消息的三种方式
public class SendMessageWays {
private KafkaProducer<String, String> producer;
public SendMessageWays() {
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.ACKS_CONFIG, "all");
// 重试次数
props.put(ProducerConfig.RETRIES_CONFIG, 3);
// 批处理大小
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384);
// 等待时间
props.put(ProducerConfig.LINGER_MS_CONFIG, 1);
producer = new KafkaProducer<>(props);
}
// 方式1:发送并忘记(不关心结果)
public void sendAndForget() {
ProducerRecord<String, String> record =
new ProducerRecord<>("my-topic", "message1");
producer.send(record);
}
// 方式2:同步发送(等待结果)
public void sendSync() throws Exception {
ProducerRecord<String, String> record =
new ProducerRecord<>("my-topic", "message2");
try {
// send()返回Future,调用get()等待完成
RecordMetadata metadata = producer.send(record).get();
System.out.printf("消息发送成功: topic=%s, partition=%d, offset=%d%n",
metadata.topic(), metadata.partition(), metadata.offset());
} catch (Exception e) {
System.err.println("消息发送失败: " + e.getMessage());
throw e;
}
}
// 方式3:异步发送(带回调)
public void sendAsync() {
ProducerRecord<String, String> record =
new ProducerRecord<>("my-topic", "message3");
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception exception) {
if (exception != null) {
System.err.println("消息发送失败: " + exception.getMessage());
} else {
System.out.printf("消息发送成功: topic=%s, partition=%d, offset=%d%n",
metadata.topic(), metadata.partition(), metadata.offset());
}
}
});
// Java 8 Lambda写法
producer.send(record, (metadata, exception) -> {
if (exception != null) {
System.err.println("消息发送失败: " + exception.getMessage());
} else {
System.out.println("消息发送成功,分区: " + metadata.partition());
}
});
}
}
带重试和异常处理的生产者
import java.util.concurrent.TimeUnit;
public class RobustProducer {
private KafkaProducer<String, String> producer;
public RobustProducer() {
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.ACKS_CONFIG, "all"); // 等待所有副本确认
props.put(ProducerConfig.RETRIES_CONFIG, 3); // 重试次数
props.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 1000); // 重试间隔
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 幂等性
// 性能配置
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "snappy"); // 压缩
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); // 32MB缓冲区
producer = new KafkaProducer<>(props);
}
public void sendWithRetry(String topic, String message) {
ProducerRecord<String, String> record =
new ProducerRecord<>(topic, message);
int maxRetries = 3;
int retryCount = 0;
while (retryCount < maxRetries) {
try {
RecordMetadata metadata = producer.send(record).get(5, TimeUnit.SECONDS);
System.out.println("消息发送成功");
return;
} catch (Exception e) {
retryCount++;
System.err.printf("发送失败(第%d次重试): %s%n", retryCount, e.getMessage());
if (retryCount >= maxRetries) {
System.err.println("达到最大重试次数,消息发送失败");
// 可以记录到死信队列或日志
handleDeadLetter(topic, message, e);
}
try {
Thread.sleep(1000 * retryCount);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
break;
}
}
}
}
private void handleDeadLetter(String topic, String message, Exception e) {
// 处理死信消息的逻辑
System.err.printf("死信消息: topic=%s, message=%s, error=%s%n",
topic, message, e.getMessage());
}
public void close() {
if (producer != null) {
producer.close();
}
}
}
批量发送消息
public class BatchProducer {
private KafkaProducer<String, String> producer;
public BatchProducer() {
Properties props = new Properties();
// ... 配置同上
// 批量配置
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 32768); // 32KB
props.put(ProducerConfig.LINGER_MS_CONFIG, 100); // 等待100ms
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 67108864); // 64MB
producer = new KafkaProducer<>(props);
}
public void sendBatchMessages(String topic, List<String> messages) {
List<ProducerRecord<String, String>> records = new ArrayList<>();
// 创建批量记录
for (int i = 0; i < messages.size(); i++) {
ProducerRecord<String, String> record =
new ProducerRecord<>(topic,
String.valueOf(i), // key
messages.get(i)); // value
records.add(record);
}
// 批量发送
List<Future<RecordMetadata>> futures = new ArrayList<>();
for (ProducerRecord<String, String> record : records) {
futures.add(producer.send(record));
}
// 等待所有发送完成
for (Future<RecordMetadata> future : futures) {
try {
RecordMetadata metadata = future.get();
System.out.printf("消息发送到分区 %d,偏移量 %d%n",
metadata.partition(), metadata.offset());
} catch (Exception e) {
System.err.println("批量发送失败: " + e.getMessage());
}
}
}
}
Maven依赖
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.5.0</version>
</dependency>
关键配置说明
| 配置项 | 默认值 | 说明 |
|---|---|---|
acks |
1 |
0:不等待确认;1:等待leader确认;all:等待所有副本确认 |
retries |
0 |
发送失败重试次数 |
batch.size |
16384 |
批次大小(字节),适当增大可以提高吞吐量 |
linger.ms |
0 |
等待更多消息加入批次的时间(毫秒) |
compression.type |
none |
压缩类型:gzip、snappy、lz4、zstd |
最佳实践
- 合理设置acks:对可靠性要求高设置
all,可以容忍丢失设置1 - 启用幂等性:设置
enable.idempotence=true防止重复消息 - 使用异步发送:配合回调函数,避免阻塞主线程
- 批量发送:适当增大batch.size和linger.ms提高性能
- 异常处理:实现重试机制和死信处理
- 资源管理:使用完务必关闭producer释放资源
这样,你就可以根据业务需求选择合适的发送方式了。