
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 groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.5.0/version /dependency如果是 Gradle 项目implementation org.apache.kafka:kafka-clients:3.5.03. 生产者配置详解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.serversbroker地址列表生产环境建议配置多个acks消息确认机制all(最安全)retries发送失败重试次数3-5batch.size批量发送大小16384-65536linger.ms发送等待时间5-1003.2 序列化器选择Kafka 消息的 key 和 value 都需要指定序列化器。除了内置的 StringSerializer还可以使用ByteArraySerializerIntegerSerializerJSON 序列化如 JacksonAvro 序列化对于复杂对象推荐使用 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 基础发送模式创建生产者并发送消息的基本流程KafkaProducerString, String producer new KafkaProducer(props); try { for(int i 0; i 100; i) { ProducerRecordString, 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); KafkaProducerString, 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(); ProducerRecordString, 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.serverslocalhost:9092 \ batch.size16384 \ linger.ms09.2 优化方向根据测试结果可能的优化点增加批量大小batch.size调整等待时间linger.ms启用压缩compression.type增加生产者实例数优化网络配置10. 与其他系统集成10.1 Spring Kafka 集成Spring Boot 提供了便捷的 Kafka 集成Configuration public class KafkaConfig { Bean public ProducerFactoryString, String producerFactory() { MapString, 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 KafkaTemplateString, String kafkaTemplate() { return new KafkaTemplate(producerFactory()); } } Service public class MessageService { Autowired private KafkaTemplateString, String kafkaTemplate; public void send(String topic, String message) { kafkaTemplate.send(topic, message); } }10.2 与流处理系统集成Kafka 消息可以被 Flink、Spark Streaming 等系统消费// Flink 消费 Kafka 示例 FlinkKafkaConsumerString consumer new FlinkKafkaConsumer( input-topic, new SimpleStringSchema(), properties); DataStreamString stream env.addSource(consumer);