Java生产者案例怎么发送消息

wen java案例 28

本文目录导读:

Java生产者案例怎么发送消息

  1. 基础配置和生产消息
  2. 发送消息的三种方式
  3. 带重试和异常处理的生产者
  4. 批量发送消息
  5. Maven依赖
  6. 关键配置说明
  7. 最佳实践

我来详细介绍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 压缩类型:gzipsnappylz4zstd

最佳实践

  1. 合理设置acks:对可靠性要求高设置all,可以容忍丢失设置1
  2. 启用幂等性:设置enable.idempotence=true防止重复消息
  3. 使用异步发送:配合回调函数,避免阻塞主线程
  4. 批量发送:适当增大batch.size和linger.ms提高性能
  5. 异常处理:实现重试机制和死信处理
  6. 资源管理:使用完务必关闭producer释放资源

这样,你就可以根据业务需求选择合适的发送方式了。

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