小菜鸟

java菜鸟号正在起航

Kafka 生产者 API 详解与实践

Kafka 生产者(Producer)是消息的发送端,负责将业务数据发送到 Kafka 集群。通过 Kafka 提供的 Producer API,我们可以灵活配置生产者行为,支持同步 / 异步发送、自定义分区策略、消息压缩等功能。本文将详细介绍生产者的实现步骤、核心配置及不同发送方式的代码示例

生产者核心概念与配置

核心配置参数

创建生产者时,需通过 Properties 对象配置关键参数,其中必选参数有三个:

参数名 作用 示例值
bootstrap.servers 指定 Kafka 集群地址(多个用逗号分隔) localhost:9092
key.serializer 消息键(Key)的序列化类(需实现 Serializer 接口) org.apache.kafka.common.serialization.StringSerializer
value.serializer 消息值(Value)的序列化类 org.apache.kafka.common.serialization.StringSerializer

常用可选参数(参考 ProducerConfig 类):

参数名 作用 默认值
acks 消息确认级别(0:不确认;1:Leader 确认;-1/all:Leader + 所有 ISR 副本确认) 1
retries 发送失败后的重试次数 0
batch.size 批次大小(达到该值后批量发送,单位:字节) 16384(16KB)
linger.ms 批处理等待时间(若未达 batch.size,超时后也会发送) 0(立即发送)
buffer.memory 发送缓冲区大小(消息暂存此处等待发送) 33554432(32MB)
compression.type 消息压缩算法(none/gzip/snappy/lz4) none

生产者工作原理

Kafka 生产者发送消息的流程如下:

阅读全文 »

Spring 集成 Kafka 详解:生产者、消费者与监听机制

Spring Kafka 是 Spring 生态对 Kafka 的集成封装,简化了 Kafka 生产者和消费者的配置与使用。它通过注解驱动、模板类(KafkaTemplate)和监听容器(MessageListenerContainer),实现了与 Spring 应用的无缝对接。本文将详细介绍 Spring Kafka 的核心组件、生产者 / 消费者实现及监听器生命周期管理。

Spring Kafka 核心组件

Spring Kafka 的设计遵循 Spring 一贯的 “模板 + 注解” 风格,核心组件包括:

组件 作用描述
KafkaTemplate 生产者模板类,封装了 Kafka 生产者 API,提供同步 / 异步发送消息的方法。
@KafkaListener 消费者注解,标注在方法上即可监听指定主题,自动接收并处理消息。
MessageListenerContainer 消息监听容器,管理消费者线程和生命周期,有单线程(KafkaMessageListenerContainer)和多线程(ConcurrentMessageListenerContainer)两种实现。
KafkaListenerEndpointRegistry 监听器容器的注册表,用于控制监听器的启动、暂停、恢复等生命周期操作。
ConsumerFactory/ProducerFactory 消费者 / 生产者工厂,负责创建 Kafka 消费者 / 生产者实例,封装配置参数。

环境配置

引入依赖

在 Maven 项目中添加 spring-kafka 依赖(需与 Kafka 版本兼容,如 Kafka 2.3.x 对应 Spring Kafka 2.3.x):

阅读全文 »

Redis 安装指南:从 macOS 到源码编译的完整步骤

Redis(Remote Dictionary Server)是一款高性能的开源键值对数据库,支持多种数据结构,广泛用于缓存、会话存储、消息队列等场景。本文详细介绍在 macOS 系统中通过包管理工具(Homebrew)和源码编译两种方式安装 Redis 的步骤,以及基本的服务管理操作。

macOS 下通过 Homebrew 安装(推荐)

Homebrew 是 macOS 下的包管理工具,通过它安装 Redis 简单高效,适合大多数用户。

1. 安装 Homebrew(若未安装)

打开终端,执行以下命令安装 Homebrew:

/bin/bash -c "$(curl -fsSL https://raw.githubusercontent.com/Homebrew/install/HEAD/install.sh)"

2. 安装 Redis

brew install redis
  • 该命令会自动下载并安装最新稳定版 Redis,同时配置环境变量,确保 redis-serverredis-cli 等命令可直接使用。

3. 配置 Redis(可选)

Redis 的配置文件 redis.conf 是核心配置入口,通过 Homebrew 安装的 Redis 配置文件路径可通过以下命令查看:

brew list redis  # 列出 Redis 安装的所有文件,其中包含 redis.conf 的路径(通常为 /usr/local/etc/redis.conf)

常用配置修改(使用文本编辑器打开 redis.conf):

  • 后台启动:默认 Redis 以前台模式运行,修改为后台启动:

阅读全文 »

Java NIO 详解:非阻塞 IO 与多路复用技术

Java NIO(Non-blocking IO,非阻塞 IO)是 JDK 1.4 引入的全新 IO 模型(JDK 1.7 补充 NIO.2),旨在解决传统 IO(BIO)在高并发场景下的性能瓶颈。NIO 基于 “通道(Channel)” 和 “缓冲区(Buffer)” 实现,通过 “多路复用(Selector)” 机制支持单线程处理多个 IO 操作,显著提升系统吞吐量。本文将全面解析 NIO 的核心原理、组件及实践。

NIO 核心概念与优势

阻塞 vs 非阻塞

  • 阻塞 IO(BIO):线程调用 read()write() 时会被挂起,直到操作完成才能继续执行。为处理多个客户端,需为每个连接创建独立线程,导致线程资源耗尽(上下文切换开销大)。
  • 非阻塞 IO(NIO):线程发起 IO 操作后无需阻塞,可继续处理其他任务;若操作未完成,仅返回 “未就绪” 状态,通过定期轮询或事件通知获取结果。单线程可管理多个 IO 通道,减少线程数量。

NIO 与传统 IO 的核心区别

特性 传统 IO(BIO) NIO
数据操作单位 字节流 / 字符流(Stream) 缓冲区(Buffer)
传输方向 单向(输入流 / 输出流分离) 双向(通道 Channel 可读写)
阻塞模式 阻塞(线程挂起) 非阻塞(线程可并发处理多任务)
并发处理 多线程(一个连接一个线程) 单线程 / 少线程(多路复用)
核心模型 流模型 通道 - 缓冲区模型
适用场景 低并发、简单 IO 操作 高并发、大流量场景(如网络服务器)

NIO 核心组件

NIO 的核心由三大组件构成:缓冲区(Buffer)通道(Channel)选择器(Selector),三者协同实现非阻塞 IO 操作。

缓冲区(Buffer):数据的容器

Buffer 是一块内存区域,用于存储 IO 操作的数据。所有 NIO 数据读写都必须通过 Buffer 完成(Channel 仅负责传输,不存储数据)。

(1)核心 Buffer 类型

NIO 为每种基本数据类型提供了对应的 Buffer 实现(除 boolean):

类型 描述 示例
ByteBuffer 字节缓冲区(最常用) 网络数据、文件二进制数据
CharBuffer 字符缓冲区 文本数据(自动处理编码)
IntBuffer/LongBuffer 基本类型缓冲区 结构化数据(如整数数组)
(2)Buffer 的核心变量

Buffer 通过三个核心变量控制数据读写,其关系为:0 ≤ mark ≤ position ≤ limit ≤ capacity

阅读全文 »

Java IO 全解析:从基础流到序列化

Java 的 I/O(输入 / 输出)机制是程序与外部世界(文件、网络、控制台等)交互的核心,其设计围绕 “流(Stream)” 展开 —— 通过流的方式有序传输数据。本文将从 I/O 操作模式出发,详细解析 Java IO 体系的核心类、使用方法及最佳实践。

I/O 操作模式:五种基本类型

I/O 操作的本质是 “数据从设备到内核缓冲区,再从缓冲区到用户进程” 的过程。根据等待数据的方式不同,分为五种模式:

模式 核心特点 适用场景
阻塞 I/O 进程发起请求后阻塞,直到数据复制完成才唤醒。 简单场景(如单线程读取小文件)
非阻塞 I/O 进程不阻塞,定期轮询缓冲区状态,数据就绪后再复制(复制时可能阻塞)。 需快速响应的场景
I/O 复用 单进程监控多个 I/O 通道,数据就绪后通知进程处理(如 select/epoll)。 高并发网络编程(如 NIO)
信号驱动 I/O 进程不阻塞,数据就绪后内核通过信号通知,再处理复制。 实时性要求高的场景
异步 I/O 进程发起请求后完全不阻塞,内核自动完成全流程,完成后通知进程。 高性能 I/O 场景(如磁盘操作)

注:Java 传统 IO(java.io)主要基于阻塞 I/O,而 NIO(java.nio)引入了 I/O 复用机制。

Java IO 核心体系:流的分类

Java IO 包(java.io)的类按功能可分为四大类,核心是 “字节流” 和 “字符流”:

类型 核心接口 / 类 处理数据类型 典型用途
字节流 InputStream(输入)、OutputStream(输出) 二进制数据(字节) 图片、视频、压缩文件等
字符流 Reader(输入)、Writer(输出) 文本数据(字符) 文本文件、配置文件等
磁盘操作 File 文件 / 目录元数据 创建 / 删除文件、获取路径
网络操作(java.net SocketServerSocket 网络字节流 客户端 / 服务器通信

字节流:处理二进制数据

字节流以字节(8 位)为单位传输数据,是所有 I/O 操作的基础。核心基类为抽象类 InputStream(输入)和 OutputStream(输出)。

字节输入流(InputStream

InputStream

所有字节输入流均继承自 InputStream,用于从数据源读取字节。

核心子类及功能
数据源 核心功能 构造器参数示例
ByteArrayInputStream 内存字节数组 从内存缓冲区读取数据,无需磁盘 I/O。 new byte[] {1,2,3}
FileInputStream 本地文件 从文件读取字节,是文件输入的基础类。 "/data/test.bin"new File(...)
PipedInputStream 管道输出流 PipedOutputStream 配合,实现线程间通信。 new PipedOutputStream()
SequenceInputStream 多个输入流 合并多个流为一个,按顺序读取。 new InputStream[] {in1, in2}
装饰器子类(FilterInputStream

通过 “装饰器模式” 为基础流添加功能(如缓冲、数据类型转换):

阅读全文 »
0%