
01
生产者应用层 (Producer Application Layer)
// 这段代码展示了生产者的核心发送流程:// 1. 消息发送入口 send() 方法// 2. 消息拦截器的调用// 3. 序列化和分区选择// 4. 消息追加到累加器// 5. 唤醒发送线程//KafkaProducer.java 的核心发送方法public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) {// 检查生产者是否已经关闭throwIfProducerClosed();// 开始记录发送时间long nowMs = time.milliseconds();// 调用拦截器ProducerRecord<K, V> interceptedRecord = this.interceptors.onSend(record);// 执行序列化和分区选择return doSend(interceptedRecord, callback, nowMs);}// 异步发送的具体实现private Future<RecordMetadata> doSend(ProducerRecord<K, V> record, Callback callback, long timestamp) {// 获取或创建topic元数据ClusterAndWaitTime clusterAndWaitTime = waitOnMetadata(record.topic(), record.partition(), maxBlockTimeMs);Cluster cluster = clusterAndWaitTime.cluster;// 序列化key和valuebyte[] serializedKey = keySerializer.serialize(record.topic(), record.headers(), record.key());byte[] serializedValue = valueSerializer.serialize(record.topic(), record.headers(), record.value());// 选择分区int partition = partition(record, serializedKey, serializedValue, cluster);// 构建消息批次并追加到累加器RecordAccumulator.RecordAppendResult result = accumulator.append(...)// 如果需要,唤醒sender线程if (result.batchIsFull || result.newBatchCreated) {this.sender.wakeup();}return result.future;}
02
消息预处理 (Message Preprocessing)
消息预处理层主要负责在消息发送前的拦截和序列化工作。ProducerInterceptors 拦截器链允许在消息发送前后进行自定义处理,比如添加时间戳、进行消息过滤、修改消息内容等。序列化器(Serializer)将用户的消息对象转换为字节数组,以便在网络上传输。Kafka 提供了多种内置的序列化器(如 String、Long 等),同时也支持自定义序列化器以满足特定的业务需求。这一层的处理对消息的格式和内容有直接影响,需要确保序列化和反序列化的一致性。
// 这部分代码展示了:// 1. 拦截器链的执行过程// 2. 序列化器的接口定义// 3. 具体序列化器的实现示例// ProducerInterceptors.javapublic ProducerRecord<K, V> onSend(ProducerRecord<K, V> record) {ProducerRecord<K, V> interceptRecord = record;for (ProducerInterceptor<K, V> interceptor : this.interceptors) {try {interceptRecord = interceptor.onSend(interceptRecord);} catch (Exception e) {// 处理异常}}return interceptRecord;}// Serializer接口public interface Serializer<T> extends Closeable {void configure(Map<String, ?> configs, boolean isKey);byte[] serialize(String topic, T data);byte[] serialize(String topic, Headers headers, T data);void close();}// StringSerializer示例public class StringSerializer implements Serializer<String> {@Overridepublic byte[] serialize(String topic, String data) {if (data == null)return null;return data.getBytes(StandardCharsets.UTF_8);}}
03
分区管理 (Partition Management)
分区管理层决定消息将被发送到主题的哪个分区。Partitioner 分区器提供了多种分区策略:轮询分区(Round-Robin)确保消息均匀分布,随机分区(Random)提供随机性,基于 key 的 Hash 分区确保相同 key 的消息总是发送到相同分区,还支持自定义分区策略以满足特定的业务需求。合理的分区策略对于负载均衡和消息顺序性有重要影响,特别是在需要保证某些消息顺序性的场景下,选择合适的分区策略尤为重要。
// 分区器代码展示了:// 1.默认分区策略的实现// 2.基于key的hash分区// 3.无key时的轮询分区// DefaultPartitioner.javapublic int partition(String topic, Object key, byte[] keyBytes,Object value, byte[] valueBytes, Cluster cluster) {List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);int numPartitions = partitions.size();if (keyBytes == null) {// 如果没有key,使用轮询策略int nextValue = nextValue(topic);List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);if (availablePartitions.size() > 0) {int part = Utils.toPositive(nextValue) % availablePartitions.size();return availablePartitions.get(part).partition();} else {// 没有可用分区时,随机选择return Utils.toPositive(nextValue) % numPartitions;}} else {// 有key时,使用hash策略return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;}}
04
消息累加器 (Record Accumulator)
// 这段代码展示了:// 消息批次的管理// 批次大小的控制// 消息追加的过程// RecordAccumulator.javapublic RecordAppendResult append(TopicPartition tp,long timestamp,byte[] key,byte[] value,Header[] headers,Callback callback,long maxTimeToBlock) {// 获取或创建消息批次ProducerBatch batch = getOrCreateBatch(tp, ...);// 追加记录到批次RecordAppendResult appendResult = batch.tryAppend(timestamp, key, value, headers, callback, time.milliseconds());// 处理追加结果if (appendResult != null) {return appendResult;}// 如果当前批次已满,创建新批次byte[] serializedKey = key;byte[] serializedValue = value;int size = Math.max(this.batchSize, AbstractRecords.estimateSizeInBytesUpperBound(magic, compression, serializedKey, serializedValue, headers));batch = createBatch(tp, size);// 再次尝试追加appendResult = batch.tryAppend(timestamp, key, value, headers, callback, time.milliseconds());return appendResult;}
05
发送线程 (Sender Thread)
// 发送线程代码展示了:// 发送循环的实现// 批次的获取和发送// 网络请求的处理// Sender.javapublic void run() {while (running) {try {runOnce(); // 执行一次发送循环} catch (Exception e) {log.error("Uncaught error in kafka producer I/O thread: ", e);}}}void runOnce() {// 获取准备发送的批次long currentTimeMs = time.milliseconds();long pollTimeout = sendProducerData(currentTimeMs);// 处理响应client.poll(pollTimeout, currentTimeMs);}private long sendProducerData(long now) {// 获取可发送的消息批次Map<Integer, List<ProducerBatch>> batches = accumulator.drain(...);// 发送消息批次for (Map.Entry<Integer, List<ProducerBatch>> entry : batches.entrySet()) {sendProducerData(entry.getKey(), entry.getValue());}}
06
网络层 (Network Layer)
// 网络层代码展示了:// 网络连接的管理// 数据发送的实现// 连接状态的维护// NetworkClient.javapublic boolean ready(Node node, long now) {// 检查连接状态if (node.isEmpty())return false;// 检查是否需要创建新连接if (connectionStates.canConnect(node.idString(), now))return true;// 检查连接是否就绪return connectionStates.isReady(node.idString(), now);}// KafkaChannel.javapublic Send write() throws IOException {if (send == null)return null;int written = send.writeTo(transportLayer);if (send.completed())send = null;return send;}
07
Kafka 集群 (Kafka Cluster)
其他副本。集群的配置和部署直接影响系统的可用性、可靠性和扩展性。// 这部分代码展示了:// 生产者与集群的交互// 元数据的获取和更新// 集群状态的检查// KafkaProducer与集群交互的相关代码private ClusterAndWaitTime waitOnMetadata(String topic, Integer partition, long maxWaitMs) {// 获取集群元数据Cluster cluster = metadata.fetch();// 检查topic是否存在if (cluster.invalidTopics().contains(topic))throw new InvalidTopicException(topic);// 更新元数据metadata.add(topic);// 等待元数据更新cluster = metadata.awaitUpdate(version, maxWaitMs);return new ClusterAndWaitTime(cluster, remainingWaitMs);}
Kafka 生产者的架构设计体现了分布式系统的精髓,通过分层设计实现了高内聚低耦合的架构特点。从消息生成到最终存储,每一层都可以独立优化和扩展。生产者的性能优化策略(如批量发送、异步处理、压缩传输等)和可靠性保证机制(如消息确认、重试策略、幂等性等)相互平衡,为不同场景提供了灵活的配置选项。通过源码分析,我们可以看到 Kafka 在工程实现上的精妙之处,比如通过 RecordAccumulator 实现高效的批次管理,通过 Sender 线程实现异步发送,通过网络层的精心设计保证了高吞吐量。这些设计思想和实现方式,为我们构建高性能、可靠的分布式消息系统提供了宝贵的参考。
07
加群请添加作者

08
获取文档代码资料

推荐阅读系列文章
建议收藏 | Dinky系列总结篇 建议收藏 | Flink系列总结篇 建议收藏 | Flink CDC 系列总结篇 建议收藏 | Doris实战文章合集 建议收藏 | Seatunnel 实战文章系列合集 建议收藏 | 实时离线输数仓(数据湖)总结篇 建议收藏 | 实时离线数仓实战第一阶段总结
如果喜欢 请点个在看分享给身边的朋友




