跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 4 / 10 查看系列简介 →

Kafka 3.9 教程(三):配置详解

文章目录
  1. 目录
  2. 3.1 Broker 配置
  3. 3.1.1 基础配置
  4. 3.1.2 网络配置
  5. 3.1.3 数据保留配置
  6. 3.1.4 ZooKeeper 配置(传统模式)
  7. 3.1.5 KRaft 配置(新模式)
  8. 3.1.6 完整 Broker 配置示例
  9. 3.2 主题级配置
  10. 3.2.1 创建主题时配置
  11. 3.2.2 修改主题配置
  12. 3.2.3 完整配置项详解
  13. 3.2.4 分区配置
  14. 3.2.5 分区分配策略
  15. 3.3 生产者配置
  16. 3.3.1 必填配置
  17. 3.3.2 可靠性配置
  18. 3.3.3 性能优化配置
  19. 3.3.4 连接配置
  20. 3.3.5 完整生产者配置示例
  21. 3.3.6 生产者配置项详解
  22. 3.3.7 生产者配置最佳实践
  23. 3.4 消费者配置
  24. 3.4.1 必填配置
  25. 3.4.2 偏移量管理配置
  26. 3.4.3 消费控制配置
  27. 3.4.4 完整消费者配置示例
  28. 3.4.5 消费者配置项详解
  29. 3.4.6 消费者配置最佳实践
  30. 3.5 Kafka Connect 配置
  31. 3.5.1 通用配置
  32. 3.5.2 独立模式配置
  33. 3.5.3 分布式模式配置
  34. 3.6 Kafka Streams 配置
  35. 3.6.1 应用配置
  36. 3.6.2 处理配置
  37. 3.6.3 并行度配置
  38. 3.6.4 内存配置
  39. 3.6.5 完整 Streams 配置示例
  40. 附录:配置优先级
  41. 下一步学习

📌 本文是「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.mslong604800000 (7天)[-1,…]消息最大保留时间,-1表示无限制
retention.byteslong-1-1,…每个分区的最大保留大小,-1表示无限制
local.retention.mslong-2[-2,…]本地日志段保留时间,-2表示使用retention.ms
local.retention.byteslong-2[-2,…]本地日志段保留大小,-2表示使用retention.bytes
delete.retention.mslong86400000 (1天)[0,…]删除标记(Tombstone)的保留时间
file.delete.delay.mslong60000 (1分钟)[0,…]文件删除前的等待时间

日志段配置

配置项类型默认值有效值说明
segment.bytesint1073741824 (1GB)[14,…]日志段文件大小
segment.mslong604800000 (7天)[1,…]日志段滚动时间
segment.jitter.mslong0[0,…]滚动时间的随机抖动,避免同时滚动
segment.index.bytesint10485760 (10MB)[4,…]索引文件大小上限
index.interval.bytesint4096 (4KB)[0,…]索引条目间隔
preallocatebooleanfalsetrue/false是否预分配日志段文件

压缩配置

配置项类型默认值有效值说明
cleanup.policystringdeletecompact, delete清理策略,可组合使用如 “delete,compact”
compression.typestringproduceruncompressed, zstd, lz4, snappy, gzip, producer最终压缩类型
compression.gzip.levelint-1[1,…,9] 或 -1GZIP压缩级别,-1使用默认级别
compression.lz4.levelint9[1,…,17]LZ4压缩级别
compression.zstd.levelint3[-131072,…,22]ZSTD压缩级别
min.cleanable.dirty.ratiodouble0.5[0,…,1]最小可清理脏比例
min.compaction.lag.mslong0[0,…]消息最小压缩延迟
max.compaction.lag.mslong9223372036854775807[1,…]消息最大压缩延迟

消息配置

配置项类型默认值有效值说明
max.message.bytesint1048588[0,…]单条消息最大大小(压缩后)
message.max.bytesint1048588[0,…]Broker级别的最大消息大小
message.timestamp.typestringCreateTimeCreateTime, LogAppendTime消息时间戳类型
message.timestamp.before.max.mslong9223372036854775807[0,…]消息时间戳可早于Broker时间的最大值
message.timestamp.after.max.mslong9223372036854775807[0,…]消息时间戳可晚于Broker时间的最大值
message.downconversion.enablebooleantruetrue/false是否启用消息格式降级转换

刷盘配置

配置项类型默认值有效值说明
flush.messageslong9223372036854775807[1,…]强制刷盘的消息间隔
flush.mslong9223372036854775807[0,…]强制刷盘的时间间隔

副本配置

配置项类型默认值说明
min.insync.replicasint1最小同步副本数(acks=all时)
unclean.leader.election.enablebooleanfalse是否允许非ISR副本选举为Leader
leader.replication.throttled.replicasstring""需要限制复制速度的副本列表
follower.replication.throttled.replicasstring""需要限制复制速度的Follower副本列表

分层存储配置

配置项类型默认值说明
remote.storage.enablebooleanfalse是否启用分层存储
remote.log.copy.disablebooleanfalse禁用分层存储数据上传
remote.log.delete.on.disablebooleanfalse禁用分层存储时删除远程数据

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.serializerclass无实现Deserializer接口的类键序列化器(必填)
value.serializerclass无实现Deserializer接口的类值序列化器(必填)
bootstrap.serverslist""host1:port1,…Broker地址列表(必填)
buffer.memorylong33554432[0,…]生产者可用于缓冲的内存总字节数
compression.typestringnonenone, gzip, snappy, lz4, zstd压缩类型,对完整批次进行压缩
retriesint2147483647[0,…]重试次数
ssl.key.passwordpasswordnull-密钥库私钥密码
ssl.keystore.certificate.chainpasswordnull-证书链(PEM格式)
ssl.keystore.keypasswordnull-私钥(PEM格式PKCS#8)
ssl.keystore.locationstringnull-密钥库文件位置
ssl.keystore.passwordpasswordnull-密钥库密码
ssl.truststore.certificatespasswordnull-信任证书(PEM格式)
ssl.truststore.locationstringnull-信任库文件位置
ssl.truststore.passwordpasswordnull-信任库密码

批次与性能配置(中优先级)

配置项类型默认值有效值说明
batch.sizeint16384[0,…]批次大小(字节),控制默认批次大小
client.dns.lookupstringuse_all_dns_ipsuse_all_dns_ips, resolve_canonical_bootstrap_servers_onlyDNS查找方式
client.idstring""-客户端标识符
compression.gzip.levelint-1[1,…,9] 或 -1GZIP压缩级别
compression.lz4.levelint9[1,…,17]LZ4压缩级别
compression.zstd.levelint3[-131072,…,22]ZSTD压缩级别
connections.max.idle.mslong540000-关闭空闲连接的超时时间
delivery.timeout.msint120000[0,…]报告成功或失败的时间上限
linger.mslong0[0,…]发送延迟,允许累积更多记录形成批次
max.block.mslong60000[0,…]send()等方法的最大阻塞时间
max.request.sizeint1048576[0,…]请求最大大小(字节)
partitioner.classclassnull-分区选择器类
partitioner.ignore.keysbooleanfalse-是否忽略记录键进行分区

网络与安全配置(中优先级)

配置项类型默认值有效值说明
receive.buffer.bytesint32768[-1,…]TCP接收缓冲区大小
request.timeout.msint30000[0,…]请求等待响应的最大时间
sasl.client.callback.handler.classclassnull-SASL客户端回调处理器
sasl.jaas.configpasswordnull-JAAS登录上下文参数
sasl.kerberos.service.namestringnull-Kafka运行的Kerberos主体名称
sasl.login.callback.handler.classclassnull-SASL登录回调处理器
sasl.login.classclassnull-实现Login接口的类
sasl.mechanismstringGSSAPI-SASL机制
sasl.oauthbearer.jwks.endpoint.urlstringnull-OAuth JWKS端点URL
sasl.oauthbearer.token.endpoint.urlstringnull-OAuth令牌端点URL
security.protocolstringPLAINTEXTPLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL与Broker通信的协议
send.buffer.bytesint131072[-1,…]TCP发送缓冲区大小
socket.connection.setup.timeout.max.mslong30000-套接字连接建立最大时间
socket.connection.setup.timeout.mslong10000-套接字连接建立超时时间
ssl.enabled.protocolslistTLSv1.2,TLSv1.3-SSL启用的协议列表
ssl.keystore.typestringJKSJKS, PKCS12, PEM密钥库文件格式
ssl.protocolstringTLSv1.3-生成SSLContext的SSL协议
ssl.providerstringnull-SSL连接的安全提供者
ssl.truststore.typestringJKSJKS, PKCS12, PEM信任库文件格式

可靠性与事务配置(低优先级)

配置项类型默认值说明
acksstringall确认机制:0=不等待, 1=leader确认, all=全部ISR确认
auto.include.jmx.reporterbooleantrue是否自动包含JmxReporter(已弃用)
enable.idempotencebooleantrue是否启用幂等性
enable.metrics.pushbooleantrue是否启用客户端指标推送
interceptor.classeslist""拦截器类列表
max.in.flight.requests.per.connectionint5单个连接上未确认请求的最大数量
metadata.max.age.mslong300000强制刷新元数据的时间间隔
metadata.max.idle.mslong300000缓存空闲主题元数据的时间
metadata.recovery.strategystringnone元数据恢复策略:none或rebootstrap
metric.reporterslist""指标报告器类列表
metrics.num.samplesint2计算指标的样本数
metrics.recording.levelstringINFO指标记录级别
metrics.sample.window.mslong30000指标采样时间窗口
partitioner.adaptive.partitioning.enablebooleantrue是否启用自适应分区
partitioner.availability.timeout.mslong0分区可用性超时时间
reconnect.backoff.max.mslong1000重连最大退避时间
reconnect.backoff.mslong50重连基础退避时间
retry.backoff.max.mslong1000重试最大退避时间
retry.backoff.mslong100重试基础退避时间
transaction.timeout.msint60000事务最大超时时间
transactional.idstringnull事务ID,用于事务传递

3.3.7 生产者配置最佳实践

场景吞吐量优先可靠性优先
acks0 或 1all
retries0Integer.MAX_VALUE
batch.size6553616384
linger.ms20-1000-5
compression.typelz4lz4
enable.idempotencefalsetrue

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.deserializerclass无实现Deserializer接口的类键反序列化器(必填)
value.deserializerclass无实现Deserializer接口的类值反序列化器(必填)
bootstrap.serverslist""host1:port1,…Broker地址列表(必填)
fetch.min.bytesint1[0,…]服务器返回的最小数据量
group.idstringnull-消费者组唯一标识
group.protocolstringclassicCONSUMER, CLASSIC消费者组协议
heartbeat.interval.msint3000-心跳间隔时间
max.partition.fetch.bytesint1048576[0,…]每个分区返回的最大数据量
session.timeout.msint45000-检测客户端失败的会话超时时间
ssl.key.passwordpasswordnull-密钥库私钥密码
ssl.keystore.certificate.chainpasswordnull-证书链(PEM格式)
ssl.keystore.keypasswordnull-私钥(PEM格式PKCS#8)
ssl.keystore.locationstringnull-密钥库文件位置
ssl.keystore.passwordpasswordnull-密钥库密码
ssl.truststore.certificatespasswordnull-信任证书(PEM格式)
ssl.truststore.locationstringnull-信任库文件位置
ssl.truststore.passwordpasswordnull-信任库密码

偏移量与消费控制配置(中优先级)

配置项类型默认值有效值说明
allow.auto.create.topicsbooleantrue-是否自动创建主题
auto.offset.resetstringlatestlatest, earliest, none无初始偏移量时的重置策略
client.dns.lookupstringuse_all_dns_ipsuse_all_dns_ips, resolve_canonical_bootstrap_servers_onlyDNS查找方式
connections.max.idle.mslong540000-关闭空闲连接的超时时间
default.api.timeout.msint60000[0,…]客户端API默认超时时间
enable.auto.commitbooleantrue-是否自动提交偏移量
exclude.internal.topicsbooleantrue-是否排除内部主题
fetch.max.bytesint52428800[0,…]服务器返回的最大数据量
group.instance.idstringnull-消费者实例唯一标识(静态成员)
group.remote.assignorstringnull-服务端分配器
isolation.levelstringread_uncommittedread_committed, read_uncommitted事务隔离级别
max.poll.interval.msint300000[1,…]poll()调用之间的最大延迟
max.poll.recordsint500[1,…]单次poll()返回的最大记录数
partition.assignment.strategylistRangeAssignor, CooperativeStickyAssignor-分区分配策略

网络与安全配置(中优先级)

配置项类型默认值有效值说明
receive.buffer.bytesint65536[-1,…]TCP接收缓冲区大小
request.timeout.msint30000[0,…]请求等待响应的最大时间
sasl.client.callback.handler.classclassnull-SASL客户端回调处理器
sasl.jaas.configpasswordnull-JAAS登录上下文参数
sasl.kerberos.service.namestringnull-Kafka运行的Kerberos主体名称
sasl.login.callback.handler.classclassnull-SASL登录回调处理器
sasl.login.classclassnull-实现Login接口的类
sasl.mechanismstringGSSAPI-SASL机制
sasl.oauthbearer.jwks.endpoint.urlstringnull-OAuth JWKS端点URL
sasl.oauthbearer.token.endpoint.urlstringnull-OAuth令牌端点URL
security.protocolstringPLAINTEXTPLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL与Broker通信的协议
send.buffer.bytesint131072[-1,…]TCP发送缓冲区大小
socket.connection.setup.timeout.max.mslong30000-套接字连接建立最大时间
socket.connection.setup.timeout.mslong10000-套接字连接建立超时时间
ssl.enabled.protocolslistTLSv1.2,TLSv1.3-SSL启用的协议列表
ssl.keystore.typestringJKSJKS, PKCS12, PEM密钥库文件格式
ssl.protocolstringTLSv1.3-生成SSLContext的SSL协议
ssl.providerstringnull-SSL连接的安全提供者
ssl.truststore.typestringJKSJKS, PKCS12, PEM信任库文件格式

低优先级配置

配置项类型默认值说明
auto.commit.interval.msint5000自动提交偏移量的频率
auto.include.jmx.reporterbooleantrue是否自动包含JmxReporter(已弃用)
check.crcsbooleantrue是否自动检查CRC32
client.idstring""客户端标识符
client.rackstring""机架标识符
enable.metrics.pushbooleantrue是否启用客户端指标推送
fetch.max.wait.msint500服务器阻塞等待满足fetch.min.bytes的最大时间
interceptor.classeslist""拦截器类列表
metadata.max.age.mslong300000强制刷新元数据的时间间隔
metadata.recovery.strategystringnone元数据恢复策略:none或rebootstrap
metric.reporterslist""指标报告器类列表
metrics.num.samplesint2计算指标的样本数
metrics.recording.levelstringINFO指标记录级别
metrics.sample.window.mslong30000指标采样时间窗口
reconnect.backoff.max.mslong1000重连最大退避时间
reconnect.backoff.mslong50重连基础退避时间
retry.backoff.max.mslong1000重试最大退避时间
retry.backoff.mslong100重试基础退避时间

分区分配策略详解

策略类说明特点
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. 生产者/消费者代码配置

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(三):配置详解》