📌 本文是「Kafka 3.9 教程」系列第 三 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
3.1 Broker 配置
Broker 是 Kafka 集群中的服务节点,负责接收、存储和转发消息。
3.1.1 基础配置
# ========== 基础配置 ==========
# Broker ID,每个 Broker 唯一
broker.id=0
# 监听器配置
listeners=PLAINTEXT://0.0.0.0:9092
# 对外公告的监听器地址(客户端连接使用)
advertised.listeners=PLAINTEXT://localhost:9092
# 消息协议与监听器映射
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_SSL:SASL_SSL
# ========== 日志存储配置 ==========
# 日志目录(可配置多个,用逗号分隔)
log.dirs=/tmp/kafka-logs
# 默认分区数
num.partitions=1
# 默认副本因子
default.replication.factor=1
# 每个日志段的大小(字节)
log.segment.bytes=1073741824 # 1GB
# 日志段滚动时间(毫秒)
log.roll.ms=604800000 # 7天
3.1.2 网络配置
# ========== 网络配置 ==========
# 最大连接数
max.connections=2147483647
# 每个 IP 的最大连接数
max.connections.per.ip=2147483647
# 套接字发送缓冲区
socket.send.buffer.bytes=102400 # 100KB
# 套接字接收缓冲区
socket.receive.buffer.bytes=102400 # 100KB
# 请求最大大小
socket.request.max.bytes=104857600 # 100MB
# 连接空闲超时时间
connections.max.idle.ms=600000 # 10分钟
# 最小拉取数据量(字节)
fetch.min.bytes=1
# 等待拉取的最大时间(毫秒)
fetch.max.wait.ms=500
3.1.3 数据保留配置
# ========== 数据保留配置 ==========
# 消息保留时间(毫秒)
log.retention.hours=168 # 7天
# 消息保留大小(字节)
log.retention.bytes=-1 # 无限制
# 最小可删除日志段大小(字节)
log.retention.check.interval.ms=300000 # 5分钟
# 清理策略
log.cleanup.policy=delete
# 压缩比率阈值
log.cleaner.min.compaction.ratio=0.5
3.1.4 ZooKeeper 配置(传统模式)
# ========== ZooKeeper 配置 ==========
zookeeper.connect=localhost:2181/kafka
# ZooKeeper 连接超时
zookeeper.connection.timeout.ms=18000
# ZooKeeper 会话超时
zookeeper.session.timeout.ms=6000
# 是否启用 ZooKeeper ACL
zookeeper.set.acl=false
3.1.5 KRaft 配置(新模式)
# ========== KRaft 配置 ==========
# 进程角色(broker, controller)
process.roles=broker,controller
# 节点 ID
node.id=1
# Controller 投票者列表
controller.quorum.voters=1@localhost:9093
# Controller 监听器
controller.listener.names=CONTROLLER
# 元数据日志副本数
metadata.log.replication.factor=3
3.1.6 完整 Broker 配置示例
# server.properties 完整示例
# 基础配置
broker.id=0
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://localhost:9092
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT
process.roles=broker,controller
node.id=1
controller.quorum.voters=1@localhost:9093
controller.listener.names=CONTROLLER
# 日志配置
log.dirs=/data/kafka-logs
num.partitions=3
default.replication.factor=3
log.segment.bytes=1073741824
log.retention.hours=168
log.retention.check.interval.ms=300000
# 性能配置
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600
# 副本配置
min.insync.replicas=2
unclean.leader.election.enable=false
# 自动创建主题
auto.create.topics.enable=true
# 删除主题配置
delete.topic.enable=true
3.2 主题级配置
主题配置可以覆盖 Broker 默认值,精细控制每个主题的行为。
3.2.1 创建主题时配置
# 创建主题并指定配置
kafka-topics.sh --create \
--topic my-topic \
--bootstrap-server localhost:9092 \
--partitions 6 \
--replication-factor 3 \
--config cleanup.policy=delete \
--config retention.ms=604800000 \
--config compression.type=lz4
3.2.2 修改主题配置
# 修改主题配置
kafka-configs.sh --alter \
--topic my-topic \
--add-config retention.ms=259200000,cleanup.policy=compact \
--bootstrap-server localhost:9092
# 查看主题配置
kafka-configs.sh --describe \
--topic my-topic \
--bootstrap-server localhost:9092
# 删除主题配置覆盖
kafka-configs.sh --alter \
--topic my-topic \
--delete-config retention.ms \
--bootstrap-server localhost:9092
3.2.3 完整配置项详解
数据保留配置
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| retention.ms | long | 604800000 (7天) | [-1,…] | 消息最大保留时间,-1表示无限制 |
| retention.bytes | long | -1 | -1,… | 每个分区的最大保留大小,-1表示无限制 |
| local.retention.ms | long | -2 | [-2,…] | 本地日志段保留时间,-2表示使用retention.ms |
| local.retention.bytes | long | -2 | [-2,…] | 本地日志段保留大小,-2表示使用retention.bytes |
| delete.retention.ms | long | 86400000 (1天) | [0,…] | 删除标记(Tombstone)的保留时间 |
| file.delete.delay.ms | long | 60000 (1分钟) | [0,…] | 文件删除前的等待时间 |
日志段配置
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| segment.bytes | int | 1073741824 (1GB) | [14,…] | 日志段文件大小 |
| segment.ms | long | 604800000 (7天) | [1,…] | 日志段滚动时间 |
| segment.jitter.ms | long | 0 | [0,…] | 滚动时间的随机抖动,避免同时滚动 |
| segment.index.bytes | int | 10485760 (10MB) | [4,…] | 索引文件大小上限 |
| index.interval.bytes | int | 4096 (4KB) | [0,…] | 索引条目间隔 |
| preallocate | boolean | false | true/false | 是否预分配日志段文件 |
压缩配置
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| cleanup.policy | string | delete | compact, delete | 清理策略,可组合使用如 “delete,compact” |
| compression.type | string | producer | uncompressed, zstd, lz4, snappy, gzip, producer | 最终压缩类型 |
| compression.gzip.level | int | -1 | [1,…,9] 或 -1 | GZIP压缩级别,-1使用默认级别 |
| compression.lz4.level | int | 9 | [1,…,17] | LZ4压缩级别 |
| compression.zstd.level | int | 3 | [-131072,…,22] | ZSTD压缩级别 |
| min.cleanable.dirty.ratio | double | 0.5 | [0,…,1] | 最小可清理脏比例 |
| min.compaction.lag.ms | long | 0 | [0,…] | 消息最小压缩延迟 |
| max.compaction.lag.ms | long | 9223372036854775807 | [1,…] | 消息最大压缩延迟 |
消息配置
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| max.message.bytes | int | 1048588 | [0,…] | 单条消息最大大小(压缩后) |
| message.max.bytes | int | 1048588 | [0,…] | Broker级别的最大消息大小 |
| message.timestamp.type | string | CreateTime | CreateTime, LogAppendTime | 消息时间戳类型 |
| message.timestamp.before.max.ms | long | 9223372036854775807 | [0,…] | 消息时间戳可早于Broker时间的最大值 |
| message.timestamp.after.max.ms | long | 9223372036854775807 | [0,…] | 消息时间戳可晚于Broker时间的最大值 |
| message.downconversion.enable | boolean | true | true/false | 是否启用消息格式降级转换 |
刷盘配置
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| flush.messages | long | 9223372036854775807 | [1,…] | 强制刷盘的消息间隔 |
| flush.ms | long | 9223372036854775807 | [0,…] | 强制刷盘的时间间隔 |
副本配置
| 配置项 | 类型 | 默认值 | 说明 |
|---|
| min.insync.replicas | int | 1 | 最小同步副本数(acks=all时) |
| unclean.leader.election.enable | boolean | false | 是否允许非ISR副本选举为Leader |
| leader.replication.throttled.replicas | string | "" | 需要限制复制速度的副本列表 |
| follower.replication.throttled.replicas | string | "" | 需要限制复制速度的Follower副本列表 |
分层存储配置
| 配置项 | 类型 | 默认值 | 说明 |
|---|
| remote.storage.enable | boolean | false | 是否启用分层存储 |
| remote.log.copy.disable | boolean | false | 禁用分层存储数据上传 |
| remote.log.delete.on.disable | boolean | false | 禁用分层存储时删除远程数据 |
3.2.4 分区配置
# 增加分区数(注意:分区数只能增加,不能减少)
kafka-topics.sh --alter \
--topic my-topic \
--partitions 12 \
--bootstrap-server localhost:9092
# 查看分区分配
kafka-topics.sh --describe \
--topic my-topic \
--bootstrap-server localhost:9092
3.2.5 分区分配策略
# Broker 默认分区分配策略
partition.assignment.strategy=org.apache.kafka.clients.consumer.RangeAssignor
# 可用策略:
# - org.apache.kafka.clients.consumer.RangeAssignor
# - org.apache.kafka.clients.consumer.RoundRobinAssignor
# - org.apache.kafka.clients.consumer.StickyAssignor
3.3 生产者配置
3.3.1 必填配置
# ========== 必填配置 ==========
# Broker 列表
bootstrap.servers=localhost:9092
# 键序列化器
key.serializer=org.apache.kafka.common.serialization.StringSerializer
# 值序列化器
value.serializer=org.apache.kafka.common.serialization.StringSerializer
3.3.2 可靠性配置
# ========== 可靠性配置 ==========
# 确认机制
acks=all
# 重试次数
retries=2147483647
# 重试间隔
retry.backoff.ms=100
# 幂等性(Kafka 0.11.0+)
enable.idempotence=true
# 事务 ID(幂等/事务需要)
transactional.id=my-transactional-id
# 事务超时
transaction.timeout.ms=60000
3.3.3 性能优化配置
# ========== 性能配置 ==========
# 批量大小(字节)
batch.size=16384 # 16KB
# 发送延迟(毫秒),0 表示立即发送
linger.ms=0
# 缓冲区大小(字节)
buffer.memory=33554432 # 32MB
# 缓冲区满时的阻塞时间
max.block.ms=5000
# 单次请求最大大小
max.request.size=1048576 # 1MB
# 压缩类型
compression.type=lz4 # none, gzip, snappy, lz4, zstd
3.3.4 连接配置
# ========== 连接配置 ==========
# 连接超时
connections.max.idle.ms=540000 # 9分钟
# 请求超时
request.timeout.ms=30000
# 重试次数
retries=3
# 重试间隔
retry.backoff.ms=100
# 元数据刷新间隔
metadata.max.age.ms=300000 # 5分钟
3.3.5 完整生产者配置示例
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;
public class ProducerConfigExample {
public static Properties getConfig() {
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, Integer.MAX_VALUE);
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 性能配置
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); // 16KB
props.put(ProducerConfig.LINGER_MS_CONFIG, 5);
props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432L); // 32MB
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
// 连接配置
props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
return props;
}
}
3.3.6 生产者配置项详解
核心配置(高优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| key.serializer | class | 无 | 实现Deserializer接口的类 | 键序列化器(必填) |
| value.serializer | class | 无 | 实现Deserializer接口的类 | 值序列化器(必填) |
| bootstrap.servers | list | "" | host1:port1,… | Broker地址列表(必填) |
| buffer.memory | long | 33554432 | [0,…] | 生产者可用于缓冲的内存总字节数 |
| compression.type | string | none | none, gzip, snappy, lz4, zstd | 压缩类型,对完整批次进行压缩 |
| retries | int | 2147483647 | [0,…] | 重试次数 |
| ssl.key.password | password | null | - | 密钥库私钥密码 |
| ssl.keystore.certificate.chain | password | null | - | 证书链(PEM格式) |
| ssl.keystore.key | password | null | - | 私钥(PEM格式PKCS#8) |
| ssl.keystore.location | string | null | - | 密钥库文件位置 |
| ssl.keystore.password | password | null | - | 密钥库密码 |
| ssl.truststore.certificates | password | null | - | 信任证书(PEM格式) |
| ssl.truststore.location | string | null | - | 信任库文件位置 |
| ssl.truststore.password | password | null | - | 信任库密码 |
批次与性能配置(中优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| batch.size | int | 16384 | [0,…] | 批次大小(字节),控制默认批次大小 |
| client.dns.lookup | string | use_all_dns_ips | use_all_dns_ips, resolve_canonical_bootstrap_servers_only | DNS查找方式 |
| client.id | string | "" | - | 客户端标识符 |
| compression.gzip.level | int | -1 | [1,…,9] 或 -1 | GZIP压缩级别 |
| compression.lz4.level | int | 9 | [1,…,17] | LZ4压缩级别 |
| compression.zstd.level | int | 3 | [-131072,…,22] | ZSTD压缩级别 |
| connections.max.idle.ms | long | 540000 | - | 关闭空闲连接的超时时间 |
| delivery.timeout.ms | int | 120000 | [0,…] | 报告成功或失败的时间上限 |
| linger.ms | long | 0 | [0,…] | 发送延迟,允许累积更多记录形成批次 |
| max.block.ms | long | 60000 | [0,…] | send()等方法的最大阻塞时间 |
| max.request.size | int | 1048576 | [0,…] | 请求最大大小(字节) |
| partitioner.class | class | null | - | 分区选择器类 |
| partitioner.ignore.keys | boolean | false | - | 是否忽略记录键进行分区 |
网络与安全配置(中优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| receive.buffer.bytes | int | 32768 | [-1,…] | TCP接收缓冲区大小 |
| request.timeout.ms | int | 30000 | [0,…] | 请求等待响应的最大时间 |
| sasl.client.callback.handler.class | class | null | - | SASL客户端回调处理器 |
| sasl.jaas.config | password | null | - | JAAS登录上下文参数 |
| sasl.kerberos.service.name | string | null | - | Kafka运行的Kerberos主体名称 |
| sasl.login.callback.handler.class | class | null | - | SASL登录回调处理器 |
| sasl.login.class | class | null | - | 实现Login接口的类 |
| sasl.mechanism | string | GSSAPI | - | SASL机制 |
| sasl.oauthbearer.jwks.endpoint.url | string | null | - | OAuth JWKS端点URL |
| sasl.oauthbearer.token.endpoint.url | string | null | - | OAuth令牌端点URL |
| security.protocol | string | PLAINTEXT | PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL | 与Broker通信的协议 |
| send.buffer.bytes | int | 131072 | [-1,…] | TCP发送缓冲区大小 |
| socket.connection.setup.timeout.max.ms | long | 30000 | - | 套接字连接建立最大时间 |
| socket.connection.setup.timeout.ms | long | 10000 | - | 套接字连接建立超时时间 |
| ssl.enabled.protocols | list | TLSv1.2,TLSv1.3 | - | SSL启用的协议列表 |
| ssl.keystore.type | string | JKS | JKS, PKCS12, PEM | 密钥库文件格式 |
| ssl.protocol | string | TLSv1.3 | - | 生成SSLContext的SSL协议 |
| ssl.provider | string | null | - | SSL连接的安全提供者 |
| ssl.truststore.type | string | JKS | JKS, PKCS12, PEM | 信任库文件格式 |
可靠性与事务配置(低优先级)
| 配置项 | 类型 | 默认值 | 说明 |
|---|
| acks | string | all | 确认机制:0=不等待, 1=leader确认, all=全部ISR确认 |
| auto.include.jmx.reporter | boolean | true | 是否自动包含JmxReporter(已弃用) |
| enable.idempotence | boolean | true | 是否启用幂等性 |
| enable.metrics.push | boolean | true | 是否启用客户端指标推送 |
| interceptor.classes | list | "" | 拦截器类列表 |
| max.in.flight.requests.per.connection | int | 5 | 单个连接上未确认请求的最大数量 |
| metadata.max.age.ms | long | 300000 | 强制刷新元数据的时间间隔 |
| metadata.max.idle.ms | long | 300000 | 缓存空闲主题元数据的时间 |
| metadata.recovery.strategy | string | none | 元数据恢复策略:none或rebootstrap |
| metric.reporters | list | "" | 指标报告器类列表 |
| metrics.num.samples | int | 2 | 计算指标的样本数 |
| metrics.recording.level | string | INFO | 指标记录级别 |
| metrics.sample.window.ms | long | 30000 | 指标采样时间窗口 |
| partitioner.adaptive.partitioning.enable | boolean | true | 是否启用自适应分区 |
| partitioner.availability.timeout.ms | long | 0 | 分区可用性超时时间 |
| reconnect.backoff.max.ms | long | 1000 | 重连最大退避时间 |
| reconnect.backoff.ms | long | 50 | 重连基础退避时间 |
| retry.backoff.max.ms | long | 1000 | 重试最大退避时间 |
| retry.backoff.ms | long | 100 | 重试基础退避时间 |
| transaction.timeout.ms | int | 60000 | 事务最大超时时间 |
| transactional.id | string | null | 事务ID,用于事务传递 |
3.3.7 生产者配置最佳实践
| 场景 | 吞吐量优先 | 可靠性优先 |
|---|
| acks | 0 或 1 | all |
| retries | 0 | Integer.MAX_VALUE |
| batch.size | 65536 | 16384 |
| linger.ms | 20-100 | 0-5 |
| compression.type | lz4 | lz4 |
| enable.idempotence | false | true |
3.4 消费者配置
3.4.1 必填配置
# ========== 必填配置 ==========
# Broker 列表
bootstrap.servers=localhost:9092
# 消费者组 ID
group.id=my-consumer-group
# 键反序列化器
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
# 值反序列化器
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
3.4.2 偏移量管理配置
# ========== 偏移量管理 ==========
# 无初始偏移量时的策略
auto.offset.reset=earliest # earliest, latest, none
# 是否自动提交偏移量
enable.auto.commit=true
# 自动提交间隔
auto.commit.interval.ms=5000
# 是否允许消费者手动分配分区(禁用消费者组)
partition.assignment.strategy=org.apache.kafka.clients.consumer.RangeAssignor
3.4.3 消费控制配置
# ========== 消费控制 ==========
# 每次拉取最大记录数
max.poll.records=500
# 拉取超时时间
fetch.min.bytes=1
fetch.max.wait.ms=500
# 最大拉取大小(字节)
max.partition.fetch.bytes=1048576 # 1MB
# 会话超时
session.timeout.ms=10000
# 心跳间隔
heartbeat.interval.ms=3000
# 消费者最长空闲时间
max.poll.interval.ms=300000 # 5分钟
3.4.4 完整消费者配置示例
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Properties;
public class ConsumerConfigExample {
public static Properties getConfig() {
Properties props = new Properties();
// 基础配置
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "my-consumer-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
// 偏移量管理
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 5000);
// 消费控制
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 500);
props.put(ConsumerConfig.FETCH_MIN_BYTES_CONFIG, 1);
props.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, 500);
// 会话管理
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 10000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 3000);
return props;
}
}
3.4.5 消费者配置项详解
核心配置(高优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| key.deserializer | class | 无 | 实现Deserializer接口的类 | 键反序列化器(必填) |
| value.deserializer | class | 无 | 实现Deserializer接口的类 | 值反序列化器(必填) |
| bootstrap.servers | list | "" | host1:port1,… | Broker地址列表(必填) |
| fetch.min.bytes | int | 1 | [0,…] | 服务器返回的最小数据量 |
| group.id | string | null | - | 消费者组唯一标识 |
| group.protocol | string | classic | CONSUMER, CLASSIC | 消费者组协议 |
| heartbeat.interval.ms | int | 3000 | - | 心跳间隔时间 |
| max.partition.fetch.bytes | int | 1048576 | [0,…] | 每个分区返回的最大数据量 |
| session.timeout.ms | int | 45000 | - | 检测客户端失败的会话超时时间 |
| ssl.key.password | password | null | - | 密钥库私钥密码 |
| ssl.keystore.certificate.chain | password | null | - | 证书链(PEM格式) |
| ssl.keystore.key | password | null | - | 私钥(PEM格式PKCS#8) |
| ssl.keystore.location | string | null | - | 密钥库文件位置 |
| ssl.keystore.password | password | null | - | 密钥库密码 |
| ssl.truststore.certificates | password | null | - | 信任证书(PEM格式) |
| ssl.truststore.location | string | null | - | 信任库文件位置 |
| ssl.truststore.password | password | null | - | 信任库密码 |
偏移量与消费控制配置(中优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| allow.auto.create.topics | boolean | true | - | 是否自动创建主题 |
| auto.offset.reset | string | latest | latest, earliest, none | 无初始偏移量时的重置策略 |
| client.dns.lookup | string | use_all_dns_ips | use_all_dns_ips, resolve_canonical_bootstrap_servers_only | DNS查找方式 |
| connections.max.idle.ms | long | 540000 | - | 关闭空闲连接的超时时间 |
| default.api.timeout.ms | int | 60000 | [0,…] | 客户端API默认超时时间 |
| enable.auto.commit | boolean | true | - | 是否自动提交偏移量 |
| exclude.internal.topics | boolean | true | - | 是否排除内部主题 |
| fetch.max.bytes | int | 52428800 | [0,…] | 服务器返回的最大数据量 |
| group.instance.id | string | null | - | 消费者实例唯一标识(静态成员) |
| group.remote.assignor | string | null | - | 服务端分配器 |
| isolation.level | string | read_uncommitted | read_committed, read_uncommitted | 事务隔离级别 |
| max.poll.interval.ms | int | 300000 | [1,…] | poll()调用之间的最大延迟 |
| max.poll.records | int | 500 | [1,…] | 单次poll()返回的最大记录数 |
| partition.assignment.strategy | list | RangeAssignor, CooperativeStickyAssignor | - | 分区分配策略 |
网络与安全配置(中优先级)
| 配置项 | 类型 | 默认值 | 有效值 | 说明 |
|---|
| receive.buffer.bytes | int | 65536 | [-1,…] | TCP接收缓冲区大小 |
| request.timeout.ms | int | 30000 | [0,…] | 请求等待响应的最大时间 |
| sasl.client.callback.handler.class | class | null | - | SASL客户端回调处理器 |
| sasl.jaas.config | password | null | - | JAAS登录上下文参数 |
| sasl.kerberos.service.name | string | null | - | Kafka运行的Kerberos主体名称 |
| sasl.login.callback.handler.class | class | null | - | SASL登录回调处理器 |
| sasl.login.class | class | null | - | 实现Login接口的类 |
| sasl.mechanism | string | GSSAPI | - | SASL机制 |
| sasl.oauthbearer.jwks.endpoint.url | string | null | - | OAuth JWKS端点URL |
| sasl.oauthbearer.token.endpoint.url | string | null | - | OAuth令牌端点URL |
| security.protocol | string | PLAINTEXT | PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL | 与Broker通信的协议 |
| send.buffer.bytes | int | 131072 | [-1,…] | TCP发送缓冲区大小 |
| socket.connection.setup.timeout.max.ms | long | 30000 | - | 套接字连接建立最大时间 |
| socket.connection.setup.timeout.ms | long | 10000 | - | 套接字连接建立超时时间 |
| ssl.enabled.protocols | list | TLSv1.2,TLSv1.3 | - | SSL启用的协议列表 |
| ssl.keystore.type | string | JKS | JKS, PKCS12, PEM | 密钥库文件格式 |
| ssl.protocol | string | TLSv1.3 | - | 生成SSLContext的SSL协议 |
| ssl.provider | string | null | - | SSL连接的安全提供者 |
| ssl.truststore.type | string | JKS | JKS, PKCS12, PEM | 信任库文件格式 |
低优先级配置
| 配置项 | 类型 | 默认值 | 说明 |
|---|
| auto.commit.interval.ms | int | 5000 | 自动提交偏移量的频率 |
| auto.include.jmx.reporter | boolean | true | 是否自动包含JmxReporter(已弃用) |
| check.crcs | boolean | true | 是否自动检查CRC32 |
| client.id | string | "" | 客户端标识符 |
| client.rack | string | "" | 机架标识符 |
| enable.metrics.push | boolean | true | 是否启用客户端指标推送 |
| fetch.max.wait.ms | int | 500 | 服务器阻塞等待满足fetch.min.bytes的最大时间 |
| interceptor.classes | list | "" | 拦截器类列表 |
| metadata.max.age.ms | long | 300000 | 强制刷新元数据的时间间隔 |
| metadata.recovery.strategy | string | none | 元数据恢复策略:none或rebootstrap |
| metric.reporters | list | "" | 指标报告器类列表 |
| metrics.num.samples | int | 2 | 计算指标的样本数 |
| metrics.recording.level | string | INFO | 指标记录级别 |
| metrics.sample.window.ms | long | 30000 | 指标采样时间窗口 |
| reconnect.backoff.max.ms | long | 1000 | 重连最大退避时间 |
| reconnect.backoff.ms | long | 50 | 重连基础退避时间 |
| retry.backoff.max.ms | long | 1000 | 重试最大退避时间 |
| retry.backoff.ms | long | 100 | 重试基础退避时间 |
分区分配策略详解
| 策略类 | 说明 | 特点 |
|---|
| RangeAssignor | 按主题分配分区 | 默认策略,按主题范围分配 |
| RoundRobinAssignor | 轮询分配 | 将分区轮询分配给消费者 |
| StickyAssignor | 粘性分配 | 保证分配最大化平衡,保留现有分配 |
| CooperativeStickyAssignor | 协作粘性分配 | 支持协作式再平衡的粘性分配 |
3.4.6 消费者配置最佳实践
| 场景 | 配置建议 |
|---|
| 高吞吐量 | max.poll.records=1000, fetch.max.wait.ms=500 |
| 低延迟 | max.poll.records=100, fetch.max.wait.ms=100 |
| 精确一次处理 | enable.auto.commit=false, 手动提交 |
| 消费者组稳定 | session.timeout.ms=30000 |
| 静态成员 | group.instance.id=固定值 |
3.5 Kafka Connect 配置
3.5.1 通用配置
# ========== 通用配置 ==========
# Broker 列表
bootstrap.servers=localhost:9092
# 集群名称
group.id=connect-cluster
# 任务数
tasks.max=1
# 键转换器
key.converter=org.apache.kafka.connect.json.JsonConverter
# 值转换器
value.converter=org.apache.kafka.connect.json.JsonConverter
# 转换器配置
key.converter.schemas.enable=true
value.converter.schemas.enable=true
# 偏移量存储
offset.storage.topic=connect-offsets
offset.storage.replication.factor=3
# 配置存储
config.storage.topic=connect-configs
config.storage.replication.factor=3
# 状态存储
status.storage.topic=connect-status
status.storage.replication.factor=3
3.5.2 独立模式配置
# connect-standalone.properties
# Worker 配置
bootstrap.servers=localhost:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
offset.storage.file.filename=/tmp/connect.offsets
# 插件路径
plugin.path=/opt/kafka/plugins
3.5.3 分布式模式配置
# connect-distributed.properties
# Worker 配置
bootstrap.servers=localhost:9092
group.id=connect-cluster
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
# 内部主题配置
offset.storage.topic=connect-offsets
offset.storage.replication.factor=3
config.storage.topic=connect-configs
config.storage.replication.factor=3
status.storage.topic=connect-status
status.storage.replication.factor=3
# REST API 配置
rest.port=8083
rest.advertised.port=8083
rest.advertised.host.name=connect-host
3.6 Kafka Streams 配置
3.6.1 应用配置
# ========== 应用配置 ==========
# 应用 ID(必须唯一)
application.id=my-streams-app
# Broker 列表
bootstrap.servers=localhost:9092
# 状态目录
state.dir=/tmp/kafka-streams
3.6.2 处理配置
# ========== 处理配置 ==========
# 处理保证
processing.guarantee=exactly_once_v2 # at_least_once, exactly_once, exactly_once_v2
# 提交间隔
commit.interval.ms=1000
# 拓扑超时
topology.max.executors=100
topology.optimization=all
3.6.3 并行度配置
# ========== 并行度配置 ==========
# Streams 线程数
num.stream.threads=1
# 任务超时
task.timeout.ms=300000 # 5分钟
# 分区缓冲区大小
buffered.records.per.partition=1000
3.6.4 内存配置
# ========== 内存配置 ==========
# 最大缓存大小
cache.max.bytes.buffering=10485760 # 10MB
# 缓存刷新阈值
cache.flush.percent=0.2
# RocksDB 配置
rocksdb.config.setter=com.example.MyRocksDBConfig
3.6.5 完整 Streams 配置示例
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.common.serialization.Serdes;
import java.util.Properties;
public class StreamsConfigExample {
public static Properties getConfig() {
Properties props = new Properties();
// 基础配置
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
// 序列化
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 处理保证
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);
// 提交间隔
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000);
// 状态目录
props.put(StreamsConfig.STATE_DIR_CONFIG, "/tmp/kafka-streams");
return props;
}
}
附录:配置优先级
配置生效优先级(从高到低):
1. 动态主题配置(kafka-configs.sh)
2. 主题创建时的 --config 参数
3. Broker 默认配置(server.properties)
4. 生产者/消费者代码配置
下一步学习