无消息丢失配置Java版本producer用户采用异步发送机制。KafkaProducer.send 方法仅仅把消息放入缓冲区中由一个专属 I/O 线程负责从缓冲区中提取消息并封装进消息batch中然后发送出去。显然这个过程中存在着数据丢失的窗口若 I/O线程发送之前 producer崩溃则存储缓冲区中的消息全部丢失了。producer 的另一个问题就是消息的乱序。假设客户端依次执行下面的语句发送两条消息到相同的分区producer.send(record1);producer.send(record2);若此时由于某些原因比如瞬时的网络抖动导致record1未发送成功同时 Kafka 又配置了重试机制以及max.in.flight.requests.per.connection大于1默认值是5那么producer重试 record1成功后record1在日志中的位置反而位于record2之后这样造成了消息的乱序。解决方案配置1.block.on.buffer.fulltrue2.acksall or-13.retriesInteger.MAX_VALUE4.max.in.flight.requests.per.connection1使用带回调机制的send发送消息即KafkaProducer.sendrecord,callbackCallback逻辑中显式地立即关闭producer使用close01.unclean.leader.election.enablefalse2.replication.factor33.min.insync.replicas24.replication.factormin.insync.replicas5.enable.auto.commitfalse分别从producer端和broker端解释参数的含义producer端block.on.buffer.fulltrue实际上这个参数在 Kafka 0.9.0.0版本已经被标记为“deprecated”并使用 max.block.ms参数替代但这里还是推荐用户显式地设置它为 true使得内存缓冲区被填满时 producer处于阻塞状态并停止接收新的消息而不是抛出异常否则producer生产速度过快会耗尽缓冲区。新版本Kafka0.10.0.0之后可以不用理会这个参数转而设置max.block.ms即可。acksall设置 acks为 all很容易理解即必须要等到所有 follower都响应了发送消息才能认为提交成功这是producer端最强程度的持久化保证。retriesInteger.MAX_VALUE设置成 MAX_VALUE纵然有些极端但其实想表达的是 producer要开启无限重试。用户不必担心producer会重试那些肯定无法恢复的错误当前producer只会重试那些可恢复的异常情况所以放心地设置一个比较大的值通常能很好地保证消息不丢失。max.in.flight.requests.per.connection1设置该参数为1主要是为了防止 topic 同分区下的消息乱序问题。这个参数的实际效果其实限制了producer在单个broker连接上能够发送的未响应请求的数量。因此如果设置成1则producer在某个broker发送响应之前将无法再给该broker发送PRODUCE请求。使用带有回调机制的send不要使用KafkaProducer中单参数的send方法因为该send调用仅仅是把消息发出而不会理会消息发送的结果。如果消息发送失败该方法不会得到任何通知故可能造成数据的丢失。实际环境中一定要使用带回调机制的send版本即KafkaProducer.sendrecord,callbackCallback逻辑中显式立即关闭producer在 Callback的失败处理逻辑中显式调用 KafkaProducer.close0。这样做的目的是为了处理消息的乱序问题。若不使用 close0默认情况下producer 会被允许将未完成的消息发送出去这样就有可能造成消息乱序。broker端配置unclean.leader.election.enablefalse关闭uncleanleader选举即不允许非ISR中的副本被选举为leader从而避免broker端因日志水位截断而造成的消息丢失。replication.factor 3设置成3主要是参考了 Hadoop及业界通用的三备份原则其实这里想强调的是一定要使用多个副本来保存分区的消息。min.insync.replicas 1用于控制某条消息至少被写入到ISR中的多少个副本才算成功设置成大于1是为了提升producer端发送语义的持久性。只有在producer端acks被设置成all或-1时这个参数才有意义。在实际使用时不要使用默认值。确保replication.factor min.insync.replicas若两者相等那么只要有一个副本挂掉分区就无法正常工作虽然有很高的持久性但可用性被极大地降低了。推荐配置成replication.factor min.insyn.replicas 1。5.1 核心概念5.1.1 为什么需要序列化 / 反序列化在网络中发送数据都是以字节的方式Kafka 也不例外。Apache Kafka支持用户给 broker发送各种类型的消息。它可以是一个字符串、一个整数、一个数组或是其他任意的对象类型。序列化器serializer负责在 producer发送前将消息转换成字节数组而与之相反解序列化器deserializer则用于将consumer接收到的字节数组转换成相应的对象。生产者 Serializer序列化器发送前把 Java 对象 → byte[]消费者 Deserializer反序列化器拉取消息后把 byte[] → Java 对象所有序列化实现都统一实现 Kafka 原生接口:// 序列化顶层接口publicinterfaceSerializerTextendsCloseable{// 初始化配置voidconfigure(MapString,?configs,booleanisKey);// 核心序列化方法对象转字节数组byte[]serialize(Stringtopic,Tdata);}// 反序列化顶层接口publicinterfaceDeserializerTextendsCloseable{voidconfigure(MapString,?configs,booleanisKey);// 核心反序列化字节数组转对象Tdeserialize(Stringtopic,byte[]data);}5.1.2 key/value 独立序列化Kafka的消息结构分为 key 和 value二者可以使用完全不同的序列化器spring:kafka:producer:# key序列化器 key-serializer:xxxSerializer # value序列化器 value-serializer:xxxSerializer1.key用于分区路由相同 key 哈希到同一分区保证有序一般用字符串序列化2.value业务主体消息字符串/JSON/Protobuf/Avro都可以5.1.3 Kafka 内置原生序列化器kafka-clients 自带无需额外依赖Kafka 客户端Java提供了常见类型的序列化器全类名在 org.apache.kafka.common.serialization 包下。props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());或者通过yml指定spring:kafka:producer:key-serializer:org.apache.kafka.common.serialization.StringSerializervalue-serializer:org.apache.kafka.common.serialization.StringSerializer注意不推荐Java原生序列化体积大、性能差、跨语言困难。仅作为演示更好的实践使用 JSON 序列化器spring-kafka 提供了JsonSerializer可以直接将任意对象转为JSON字节。spring:kafka:producer:value-serializer:org.springframework.kafka.support.serializer.JsonSerializer发送时直接传入 POJOkafkaTemplate.send(topic,newUser(张三,25));消费者端需配置 JsonDeserializer 并指定可信包spring:kafka:consumer:value-deserializer:org.springframework.kafka.support.serializer.JsonDeserializerproperties:spring.json.trusted.packages:*spring.json.value.default.type:com.example.User5.1.4 序列化器与反序列化器对称性1.必须成对使用生产者用什么序列化消费者就必须用对应的反序列化器。2.例如StringSerializer↔StringDeserializerJsonSerializer↔JsonDeserializer5.1.5 性能与注意事项1.吞吐量紧凑的格式Avro、Protobuf比JSON更快且体积小2.PU开销JSON序列化/反序列化中等Avro需要 schema 查找有额外网络开销SchemaRegistry原生Java序列化最差。
5. Kafka 序列化器,生产者 Serializer,消费者 Deserializer
无消息丢失配置Java版本producer用户采用异步发送机制。KafkaProducer.send 方法仅仅把消息放入缓冲区中由一个专属 I/O 线程负责从缓冲区中提取消息并封装进消息batch中然后发送出去。显然这个过程中存在着数据丢失的窗口若 I/O线程发送之前 producer崩溃则存储缓冲区中的消息全部丢失了。producer 的另一个问题就是消息的乱序。假设客户端依次执行下面的语句发送两条消息到相同的分区producer.send(record1);producer.send(record2);若此时由于某些原因比如瞬时的网络抖动导致record1未发送成功同时 Kafka 又配置了重试机制以及max.in.flight.requests.per.connection大于1默认值是5那么producer重试 record1成功后record1在日志中的位置反而位于record2之后这样造成了消息的乱序。解决方案配置1.block.on.buffer.fulltrue2.acksall or-13.retriesInteger.MAX_VALUE4.max.in.flight.requests.per.connection1使用带回调机制的send发送消息即KafkaProducer.sendrecord,callbackCallback逻辑中显式地立即关闭producer使用close01.unclean.leader.election.enablefalse2.replication.factor33.min.insync.replicas24.replication.factormin.insync.replicas5.enable.auto.commitfalse分别从producer端和broker端解释参数的含义producer端block.on.buffer.fulltrue实际上这个参数在 Kafka 0.9.0.0版本已经被标记为“deprecated”并使用 max.block.ms参数替代但这里还是推荐用户显式地设置它为 true使得内存缓冲区被填满时 producer处于阻塞状态并停止接收新的消息而不是抛出异常否则producer生产速度过快会耗尽缓冲区。新版本Kafka0.10.0.0之后可以不用理会这个参数转而设置max.block.ms即可。acksall设置 acks为 all很容易理解即必须要等到所有 follower都响应了发送消息才能认为提交成功这是producer端最强程度的持久化保证。retriesInteger.MAX_VALUE设置成 MAX_VALUE纵然有些极端但其实想表达的是 producer要开启无限重试。用户不必担心producer会重试那些肯定无法恢复的错误当前producer只会重试那些可恢复的异常情况所以放心地设置一个比较大的值通常能很好地保证消息不丢失。max.in.flight.requests.per.connection1设置该参数为1主要是为了防止 topic 同分区下的消息乱序问题。这个参数的实际效果其实限制了producer在单个broker连接上能够发送的未响应请求的数量。因此如果设置成1则producer在某个broker发送响应之前将无法再给该broker发送PRODUCE请求。使用带有回调机制的send不要使用KafkaProducer中单参数的send方法因为该send调用仅仅是把消息发出而不会理会消息发送的结果。如果消息发送失败该方法不会得到任何通知故可能造成数据的丢失。实际环境中一定要使用带回调机制的send版本即KafkaProducer.sendrecord,callbackCallback逻辑中显式立即关闭producer在 Callback的失败处理逻辑中显式调用 KafkaProducer.close0。这样做的目的是为了处理消息的乱序问题。若不使用 close0默认情况下producer 会被允许将未完成的消息发送出去这样就有可能造成消息乱序。broker端配置unclean.leader.election.enablefalse关闭uncleanleader选举即不允许非ISR中的副本被选举为leader从而避免broker端因日志水位截断而造成的消息丢失。replication.factor 3设置成3主要是参考了 Hadoop及业界通用的三备份原则其实这里想强调的是一定要使用多个副本来保存分区的消息。min.insync.replicas 1用于控制某条消息至少被写入到ISR中的多少个副本才算成功设置成大于1是为了提升producer端发送语义的持久性。只有在producer端acks被设置成all或-1时这个参数才有意义。在实际使用时不要使用默认值。确保replication.factor min.insync.replicas若两者相等那么只要有一个副本挂掉分区就无法正常工作虽然有很高的持久性但可用性被极大地降低了。推荐配置成replication.factor min.insyn.replicas 1。5.1 核心概念5.1.1 为什么需要序列化 / 反序列化在网络中发送数据都是以字节的方式Kafka 也不例外。Apache Kafka支持用户给 broker发送各种类型的消息。它可以是一个字符串、一个整数、一个数组或是其他任意的对象类型。序列化器serializer负责在 producer发送前将消息转换成字节数组而与之相反解序列化器deserializer则用于将consumer接收到的字节数组转换成相应的对象。生产者 Serializer序列化器发送前把 Java 对象 → byte[]消费者 Deserializer反序列化器拉取消息后把 byte[] → Java 对象所有序列化实现都统一实现 Kafka 原生接口:// 序列化顶层接口publicinterfaceSerializerTextendsCloseable{// 初始化配置voidconfigure(MapString,?configs,booleanisKey);// 核心序列化方法对象转字节数组byte[]serialize(Stringtopic,Tdata);}// 反序列化顶层接口publicinterfaceDeserializerTextendsCloseable{voidconfigure(MapString,?configs,booleanisKey);// 核心反序列化字节数组转对象Tdeserialize(Stringtopic,byte[]data);}5.1.2 key/value 独立序列化Kafka的消息结构分为 key 和 value二者可以使用完全不同的序列化器spring:kafka:producer:# key序列化器 key-serializer:xxxSerializer # value序列化器 value-serializer:xxxSerializer1.key用于分区路由相同 key 哈希到同一分区保证有序一般用字符串序列化2.value业务主体消息字符串/JSON/Protobuf/Avro都可以5.1.3 Kafka 内置原生序列化器kafka-clients 自带无需额外依赖Kafka 客户端Java提供了常见类型的序列化器全类名在 org.apache.kafka.common.serialization 包下。props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());或者通过yml指定spring:kafka:producer:key-serializer:org.apache.kafka.common.serialization.StringSerializervalue-serializer:org.apache.kafka.common.serialization.StringSerializer注意不推荐Java原生序列化体积大、性能差、跨语言困难。仅作为演示更好的实践使用 JSON 序列化器spring-kafka 提供了JsonSerializer可以直接将任意对象转为JSON字节。spring:kafka:producer:value-serializer:org.springframework.kafka.support.serializer.JsonSerializer发送时直接传入 POJOkafkaTemplate.send(topic,newUser(张三,25));消费者端需配置 JsonDeserializer 并指定可信包spring:kafka:consumer:value-deserializer:org.springframework.kafka.support.serializer.JsonDeserializerproperties:spring.json.trusted.packages:*spring.json.value.default.type:com.example.User5.1.4 序列化器与反序列化器对称性1.必须成对使用生产者用什么序列化消费者就必须用对应的反序列化器。2.例如StringSerializer↔StringDeserializerJsonSerializer↔JsonDeserializer5.1.5 性能与注意事项1.吞吐量紧凑的格式Avro、Protobuf比JSON更快且体积小2.PU开销JSON序列化/反序列化中等Avro需要 schema 查找有额外网络开销SchemaRegistry原生Java序列化最差。