Kafka Producer
生产者客户端由两个线程协调运行。其中主线程创建消息,并经过拦截器、序列化器、分区器作用后缓存到消息累加器; 消息累加器中的ProducerBatch是一个双端队列,消息添加时从尾部进入,Sender读取消息时从头部取出。ProducerBatch包含链多个ProducerRecord,这样使字节的使用更加紧凑,同时发送时减少了网络请求的次数以提升整体的吞吐量。
RecordAccumulator缓存的大小可以通过生产者客户端参数Buffer.memory配置,默认值为32MB。如果生产者发送消息的速度超过发送到服务器的速度,则会导致生产者空间不足,这个时候KafkaProducer的send方法要么被阻塞,要么抛出异常,这个取决于max.block.ms的配置,默认为60000,即60秒。如果生产者客户端要向很多分区发送消息,则可以适当增加buffer.memory,以提高整体的吞吐量。
在Sender线程将请求发送kafka之前,还会将请求保存到InFlightRequests中,其形式是Map
leastLoadedNode是所有Node中负载最小的那一个,选择这个Node发送请求可以使请求尽快发出,避免因网络拥塞等异常而影响发送的进度。通常用于元数据请求、消费者组播协议的交互。
在向一个主题发送消息时,KafkaProducer要将此消息追加到指定主题的某个分区所对应的leader副本之前,首先需要知道分区的数量,然后经过计算得出目标分区,还要知道目标分区的leader副本所在的broker节点的地址、端口信息之后才能建立连接,最终才可以把消息发送到kafka。
元数据是指Kafka集群的元数据,这些元数据记录链集群的有什么主题、主题有哪些分区、每个分区的leader副本分配在那个节点,follower副本分配在哪些节点上,哪些副本在AR 、ISR等集合中,集群中有哪些节点,控制器节点又是哪一个等等信息。
当客户端中没有需要使用的元数据信息时,比如没有指定的主题信息,或者超过metadata.max.age.ms 时间没有更新元数据都会引起元数据的更新操作。客户端参数metadata.max.age.ms的默认值为300000,即5分钟。元数据的更新由Sender线程向leastLoadedNode发送MetadataRequest请求。

生产者客户端架构
生产者配置参数:
1.acks
acks=1。默认值即为1。生产者发送消息之后,只要分区的leader副本成功写入消息,那么它就会收到来自服务端的成功响应。如果消息无法写入leader副本,比如在leader 副本崩溃、重新选举新的 leader 副本的过程中,那么生产者就会收到一个错误的响应,为了避免消息丢失,生产者可以选择重发消息。
acks=0。生产者发送消息之后不需要等待任何服务端的响应。acks 设置为 0 可以达到最大的吞吐量。
acks=-1或acks=all。生产者在消息发送之后,需要等待ISR中的所有副本都成功写入消息之后才能够收到来自服务端的成功响应。在其他配置环境相同的情况下,acks 设置为-1(all)可以达到最强的可靠性。
2.max.request.size
这个参数用来限制生产者客户端能发送的消息的最大值,默认值为 1048576B,即1MB。
3.retries和retry.backoff.ms
retries参数用来配置生产者重试的次数,默认值为0,即在发生异常的时候不进行任何重试动作。消息在从生产者发出到成功写入服务器之前可能发生一些临时性的异常,比如网络抖动、leader副本的选举等,这种异常往往是可以自行恢复的,生产者可以通过配置retries大于0的值,以此通过内部重试来恢复
retry.backoff.ms这个参数的默认值为100,它用来设定两次重试之间的时间间隔,避免无效的频繁重试.
4.compression.type
这个参数用来指定消息的压缩方式,默认值为“none”,即默认情况下,消息不会被压缩。该参数还可以配置为“gzip”“snappy”和“lz4”。对消息进行压缩可以极大地减少网络传输量、降低网络I/O,从而提高整体的性能。消息压缩是一种使用时间换空间的优化方式,如果对时延有一定的要求,则不推荐对消息进行压缩。
5.connections.max.idle.ms
指定在多久之后关闭限制的连接,默认值是540000(ms),即9分钟。
6.linger.ms
指定生产者发送 ProducerBatch 之前等待更多消息(ProducerRecord)加入ProducerBatch 的时间,默认值为 0。生产者客户端会在 ProducerBatch 被填满或等待时间超过linger.ms 值时发送出去。增大这个参数的值会增加消息的延迟,但是同时能提升一定的吞吐量。
7.receive.buffer.bytes
设置Socket接收消息缓冲区(SO_RECBUF)的大小,默认值为32768(B),即32KB。如果设置为-1,则使用操作系统的默认值。
8.send.buffer.bytes
设置Socket发送消息缓冲区(SO_SNDBUF)的大小,默认值为131072(B),即128KB。与receive.buffer.bytes参数一样,如果设置为-1,则使用操作系统的默认值。
9.request.timeout.ms
配置Producer等待请求响应的最长时间,默认值为30000(ms)。请求超时之后可以选择进行重试。注意这个参数需要比broker端参数replica.lag.time.max.ms的值要大,这样可以减少因客户端重试而引起的消息重复的概率。