# study-kafka **Repository Path**: zrbfree/study-kafka ## Basic Information - **Project Name**: study-kafka - **Description**: kafka3.x 基本操作 1、拦截器、序列化器、分区器的使用 2、偏移量控制 3、应答机制 4、幂等性 5、事务 - **Primary Language**: Java - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-19 - **Last Updated**: 2026-08-29 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # Kafka 学习项目 (study-kafka) 这是一个 Kafka 学习项目,通过大量的代码示例帮助理解和掌握 Kafka 的核心功能。本项目使用 Java 语言编写,基于 Spring Boot 框架。 ## 项目结构 ``` src/main/java/com/rick/ ├── StudyKafkaApplication.java # Spring Boot 启动类 ├── Constantis.java # 常量接口 ├── contracts/ │ └── MyConfig.java # Kafka 配置常量类 ├── pojo/ │ └── Student.java # 示例实体类 ├── _01_api/ │ └── KafkaTopicDML.java # Topic 管理操作 ├── _02_quickStart/ │ ├── KafkaProducerQuickStart.java # 生产者快速入门 │ ├── KafkaConsumerQuickStart_1.java # 消费者入门示例 1 │ └── KafkaConsumerQuickStart_2.java # 消费者入门示例 2 ├── _03_interceptors/ │ ├── KafkaConsumerInterceptor.java # 消费者拦截器示例 │ ├── KafkaProducerInteceptor.java # 生产者拦截器示例 │ └── UserDefineProducerInterceptor.java # 自定义生产者拦截器 ├── _04_serializer/ │ ├── KafkaProducerUser.java # 用户序列化生产者 │ ├── KafkaConsumerUser.java # 用户序列化消费者 │ ├── ProducerWithStudentSerializer.java # 学生对象序列化 │ ├── ProducerWithStudentSerializerNoClass.java │ ├── User.java # 用户实体 │ ├── UserDefineSerializer.java # 自定义序列化器 │ └── UserDefineDeserializer.java # 自定义反序列化器 ├── _05_partitionr/ │ ├── KafkaProducerAssignPartition.java # 指定分区示例 │ ├── KafkaProducerPartitioner.java # 分区器示例 │ ├── ProducerWithUDPPartitioner.java # 自定义分区器 │ └── UserDefinePartitioner.java # 自定义分区器实现 ├── _06_offset/ │ ├── KafkaConsumerOffset_1~4.java # _offset 管理系列 │ ├── KafkaProducerOffset.java # 生产者_offset │ ├── ConsumerWithAutoOffsetMysql.java # MySQL 存储_offset │ ├── ConsumerWithAutoOffsetRedis.java # Redis 存储_offset │ └── ConsumerWithUDOffset*.java # 自定义_offset 管理 ├── _07_acks/ │ ├── KafkaProducerAcks.java # 生产者_ack 配置 │ ├── KafkaConsumerAcks.java # 消费者_ack 处理 │ ├── ProducerWithCallBack.java # 带回调的生产者 │ └── ProducerWithCallBackAndACK.java # 回调+ACK 组合 ├── _08_idempotence/ │ ├── KafkaProducerIdempotence.java # 幂等性生产者 │ └── KafkaConsumerIdempotence.java # 幂等性消费者 ├── _09_transactions/ │ ├── ProducerWithTransaction.java # 事务生产者 │ ├── KafkaProducerTransaction*.java # 事务相关示例 │ └── KafkaConsumerTransaction1.java # 事务消费者 └── _10_consumer_assignor/ ├── RangeAssignorDemo.java # Range 分区分配策略 ├── RoundRobinAssignorDemo.java # RoundRobin 分配策略 ├── StickyAssignorDemo.java # Sticky 分配策略 ├── CooperativeStickyAssignor.java # CooperativeSticky 策略 └── CustomAssignorDemo.java # 自定义分配策略 ``` ## 核心功能模块 ### 1. Topic 管理 - Topic 的创建、删除、查询等 DML 操作 ### 2. 快速入门 - 生产者发送消息的基础用法 - 消费者接收消息的两种方式 ### 3. 拦截器 - 生产者拦截器:可在消息发送前、后进行拦截处理 - 消费者拦截器:可在消息消费前进行拦截处理 - 自定义拦截器实现 ### 4. 序列化与反序列化 - 自定义序列化器:实现 `Serializer` 接口 - 自定义反序列化器:实现 `Deserializer` 接口 - 支持多种数据类型的序列化 ### 5. 分区策略 - 指定分区发送 - 自定义分区器:实现 `Partitioner` 接口 - 根据 key 路由到不同分区 ### 6. Offset 管理 - 自动提交 offset - 手动提交 offset - 自定义存储 offset(MySQL、Redis) ### 7. ACK 机制 - `acks=0`、`acks=1`、`acks=all` 配置 - 回调函数处理发送结果 ### 8. 幂等性 - 开启生产者的幂等性 - 避免消息重复发送 ### 9. 事务支持 - 事务生产者:原子性发送消息 - 事务消费者:按事务消费消息 ### 10. 消费者分区分配策略 - `RangeAssignor`:范围分配 - `RoundRobinAssignor`:轮询分配 - `StickyAssignor`:粘性分配 - `CooperativeStickyAssignor`:协作式粘性分配 - 自定义分配策略 ## 配置说明 项目中的主要配置项(定义在 `MyConfig.java`): | 配置项 | 说明 | |--------|------| | `BOOTSTRAP_SERVERS` | Kafka 集群地址 | | `GROUP_ID` | 消费者组 ID | | `BATCH_SIZE` | 批量发送大小 | | `LINGER_MS` | 发送延迟时间 | | `ACK` | 确认机制 | | `RETRIES` | 重试次数 | ## 环境要求 - Java 8+ - Kafka 2.x 或更高版本 - Maven 3.x ## 快速开始 ### 1. 导入项目 使用 IDE(如 IntelliJ IDEA)导入项目,或通过 Maven 命令: ```bash mvn clean install ``` ### 2. 修改配置 在 `MyConfig.java` 中修改 Kafka 连接地址: ```java public static final String BOOTSTRAP_SERVERS = "localhost:9092"; ``` ### 3. 运行示例 直接运行对应的 `main` 方法即可,例如: - 运行 `KafkaProducerQuickStart.java` 发送消息 - 运行 `KafkaConsumerQuickStart_1.java` 接收消息 ## 学习建议 建议按照以下顺序学习: 1. **快速入门**:了解基本的发送和接收流程 2. **API 操作**:学习 Topic 的管理操作 3. **序列化/反序列化**:掌握自定义数据的编解码 4. **分区策略**:理解消息路由机制 5. **Offset 管理**:学习消费者 Offset 的管理方式 6. **ACK 与幂等**:确保消息可靠性 7. **事务**:学习跨分区、跨会话的消息原子性 8. **拦截器**:了解 Kafka 的扩展机制 ## 注意事项 - 运行前请确保 Kafka 服务已启动 - 根据实际环境修改配置文件中的连接地址 - 部分示例需要配合多个消费者实例才能看到效果 ## License 本项目仅供学习使用。