新闻详情

新闻详情

首页 / 资讯中心 / 详情

Kafka

发布时间:2026/10/1 15:41:08来源:尧图网络
Kafka
一、Kafka 环境搭建单机版1. 下载 Kafka前往 Apache Kafka 官网 下载最新稳定版例如 3.6.0。解压后目录结构如下kafka_2.13-3.6.0/ ├── bin/ # 启动脚本 ├── config/ # 配置文件 ├── libs/ # 依赖库 └── ...2. 启动 Kafka早期 Kafka 依赖 ZooKeeper新版本推荐使用KRaft模式无需 ZooKeeper。以下分别介绍两种方式。方式一KRaft 模式Kafka 3.3 推荐# 1. 生成集群 ID KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) # 2. 格式化存储目录使用默认配置 bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 3. 启动 Kafka bin/kafka-server-start.sh config/kraft/server.properties方式二ZooKeeper 模式传统# 1. 启动 ZooKeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 2. 启动 Kafka bin/kafka-server-start.sh config/server.properties默认端口Kafka:9092ZooKeeper:21813. 创建 Topicbin/kafka-topics.sh --create --topic my-topic --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1查看 Topic 列表bin/kafka-topics.sh --list --bootstrap-server localhost:90924. 命令行测试发送消息bin/kafka-console-producer.sh --topic my-topic --bootstrap-server localhost:9092 hello kafka消费消息bin/kafka-console-consumer.sh --topic my-topic --from-beginning --bootstrap-server localhost:9092二、Java 原生客户端使用在 Maven 项目中添加依赖pom.xmldependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.6.0/version /dependency1. 生产者Producer基本配置项配置项说明bootstrap.serversKafka 集群地址多个用逗号分隔key.serializer键的序列化器如StringSerializervalue.serializer值的序列化器acks确认机制0不等待确认1仅 leader 确认all或-1所有副本确认retries发送失败重试次数batch.size批量发送大小字节linger.ms等待更多消息加入批次的时间buffer.memory生产者缓冲区大小compression.type压缩类型none、gzip、snappy、lz4、zstd生产者代码示例import org.apache.kafka.clients.producer.*; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; import java.util.concurrent.Future; public class MyProducer { public static void main(String[] args) { 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); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, snappy); // 创建生产者 ProducerString, String producer new KafkaProducer(props); // 1. 发送消息异步不关心结果 producer.send(new ProducerRecord(my-topic, key1, value1)); // 2. 发送消息并获取 Future可阻塞等待结果 FutureRecordMetadata future producer.send( new ProducerRecord(my-topic, key2, value2) ); try { RecordMetadata metadata future.get(); System.out.println(发送成功offset metadata.offset() , partition metadata.partition()); } catch (Exception e) { e.printStackTrace(); } // 3. 发送消息并带回调异步 producer.send(new ProducerRecord(my-topic, key3, value3), new Callback() { Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception null) { System.out.println(发送成功: metadata.offset()); } else { exception.printStackTrace(); } } }); // 关闭生产者会等待所有缓冲消息发送完成 producer.close(); } }2. 消费者Consumer基本配置项配置项说明bootstrap.serversKafka 集群地址group.id消费者组 ID相同组内的消费者共同消费key.deserializer键的反序列化器value.deserializer值的反序列化器enable.auto.commit是否自动提交 offsetauto.offset.reset初始消费位置earliest从头开始latest从最新开始max.poll.records一次 poll 最多返回的记录数session.timeout.ms会话超时时间消费者代码示例import org.apache.kafka.clients.consumer.*; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class MyConsumer { public static void main(String[] args) { Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, my-group); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); // 自动提交 props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 1000); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest); // 从头消费 ConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(my-topic)); try { while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, String record : records) { System.out.printf(offset%d, key%s, value%s%n, record.offset(), record.key(), record.value()); } } } finally { consumer.close(); } } }手动提交 offset推荐生产环境使用props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 关闭自动提交 while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, String record : records) { // 处理消息 System.out.println(record.value()); } // 处理完后手动提交 offset同步提交 consumer.commitSync(); // 或者异步提交 // consumer.commitAsync(); }三、Spring Boot 集成 KafkaSpring Kafka 对 Kafka 进行了封装大大简化了开发。下面演示完整流程。1. 添加依赖在pom.xml中添加dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependencySpring Boot 会自动管理版本通常不需要显式指定版本号2. 配置文件application.ymlspring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 consumer: group-id: my-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest enable-auto-commit: false listener: ack-mode: manual # 手动提交 offset3. 生产者服务import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; Service public class KafkaProducerService { private final KafkaTemplateString, String kafkaTemplate; public KafkaProducerService(KafkaTemplateString, String kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } public void sendMessage(String topic, String message) { kafkaTemplate.send(topic, message); } // 带 key 和回调 public void sendMessageWithCallback(String topic, String key, String message) { kafkaTemplate.send(topic, key, message).addCallback( result - System.out.println(发送成功: result.getRecordMetadata().offset()), ex - System.err.println(发送失败: ex.getMessage()) ); } }4. 消费者监听器import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; Component public class KafkaConsumer { // 自动确认配置中 enable-auto-commit: true 时 KafkaListener(topics my-topic, groupId my-group) public void listen(String message) { System.out.println(收到消息: message); } // 手动提交 offset配置中 enable-auto-commit: false 时 KafkaListener(topics my-topic, groupId my-group) public void listenManual(String message, Acknowledgment ack) { System.out.println(收到消息: message); // 处理完成后提交 offset ack.acknowledge(); } }5. 发送消息测试在 Controller 或测试类中注入KafkaProducerService调用即可。四、进阶特性与最佳实践1. 自定义序列化器如果消息是 Java 对象可以使用 JSON 或 Avro 序列化。常用的是 Spring Kafka 提供的JsonSerializer/JsonDeserializer或者使用StringSerializer配合 Jackson 手动转换。2. 分区策略默认分区器如果指定了 key则根据 key 的 hash 选择分区如果 key 为 null则使用轮询round-robin。自定义分区器实现Partitioner接口并在配置中指定partitioner.class。3. 消息顺序性Kafka 只能保证同一个分区内消息有序。因此对于需要严格顺序的业务可以将需要有序的消息发送到同一个分区例如使用相同的 key。4. 幂等性Kafka 生产者默认开启幂等enable.idempotencetrue可以避免因网络重试导致的重复消息。对于消费者端需要业务逻辑本身具备幂等性。5. 事务Kafka 支持事务可以保证“发送多条消息要么全部成功要么全部失败”。配置props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, my-transactional-id); producer.initTransactions(); producer.beginTransaction(); producer.send(record1); producer.send(record2); producer.commitTransaction();6. 错误处理在消费者中如果处理消息失败可以抛出异常Spring Kafka 会进行重试。也可以配置ErrorHandler或SeekToCurrentErrorHandler来控制失败后的行为例如重试一定次数后跳过或发送到死信队列。五、常见问题解答Q1: 如何保证消息不丢失生产者设置acksallretries0开启幂等。Broker设置min.insync.replicas至少为 2保证至少有一个副本同步。消费者禁用自动提交处理完消息后手动提交 offset。Q2: 如何保证消息不重复Kafka 本身无法绝对避免重复需要消费者端做幂等如数据库唯一约束、Redis 去重等。可以使用事务或幂等生产者减少重复发送。Q3: 消费进度如何回溯使用kafka-consumer-groups.sh工具重置 offsetbin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --reset-offsets --to-earliest --topic my-topic --executeQ4: 如何监控 Kafka使用 Kafka 自带的 JMX 指标 Prometheus Grafana。使用 Kafka Manager、Kafka Eagle、Burrow 等第三方工具。六、总结使用 Kafka 的基本步骤启动 Kafka 集群单机或集群。创建 Topic。编写生产者配置序列化器、acks 等发送消息。编写消费者配置反序列化器、group.id订阅 Topic 并处理消息。生产环境考虑消息可靠性不丢失、不重复、顺序性、分区策略、监控等。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

货拉拉营销广告大模型落地实战:提示词工程与智能体工作流 2026/10/1 18:46:05

货拉拉营销广告大模型落地实战:提示词工程与智能体工作流

1. 货拉拉营销广告的真实痛点:为什么通用大模型直接拿来用会翻车 货拉拉的营销广告业务有个很鲜明的特点:它不是那种"一个品牌对全网喊话"的标准化投放,而是 同城货运场景下、司机端与货主端双角色、多城市多车型多时段 的碎片化…

阅读更多 →
TypeScript从入门到实践:类型系统、泛型与工程迁移指南 2026/10/1 18:46:05

TypeScript从入门到实践:类型系统、泛型与工程迁移指南

如果你写过一段时间的JavaScript,大概率经历过这种时刻:一个函数跑得好好的,换个调用方式突然就报错了;一段别人留下的老代码,改了一行数据格式,十几个地方跟着崩;又或者一个对象明明有某个字段…

阅读更多 →
Java与Python项目服务器部署实战:从环境配置到前后端分离 2026/10/1 18:46:05

Java与Python项目服务器部署实战:从环境配置到前后端分离

干开发这些年,最常见的场景就是:代码写得挺欢,一到“部署”这两个字就头疼。Java项目打包出个jar或者war,扔到服务器上跑不起来;Python项目本地运行没问题,换台机器一堆依赖报错。尤其是从“能运行”到“稳…

阅读更多 →
TypeScript 实战:从类型系统到渐进迁移 2026/10/1 18:46:04

TypeScript 实战:从类型系统到渐进迁移

1. 为什么说 TypeScript 是 JavaScript 的一次蜕变做了这么多年前端,我最初对 TypeScript 的态度也是“多此一举”。JavaScript 写得好好的,为什么要多一层编译?直到在一个中型项目里被一个undefined is not a function的报错折腾了三个小时&…

阅读更多 →
二进制与十六进制互转及float还原:大小端、移位全解析 2026/10/1 18:46:04

二进制与十六进制互转及float还原:大小端、移位全解析

前阵子帮人排查一个嵌入式设备日志,里面打了一串十六进制字节,对方问我怎么把它还原成真实的 float 数值。说实话,干这行久了,这类问题见得太多,但每次被问还是会感慨一句:二进制和十六进制,平时…

阅读更多 →
054振荡排序 2026/10/1 18:45:45

054振荡排序

振荡排序 (Oscillating Sort / Reversing Merge) 054钟摆算法:解码振荡排序故事:钟摆的节拍 在磁带机时代,有一个令工程师头疼的问题:磁带倒带很慢。每次排序合并之后,都要把磁带倒回起始位置,才能进行下一…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉