Kafka的消费全流程
我们接着继续去理解最后这条消息是如何被消费者消费掉的。其中最核心的有以下内容。1、多线程安全问题2、群组协调3、分区再均衡多线程安全问题当多个线程访问某个类时这个类始终都能表现出正确的行为那么就称这个类是线程安全的。对于线程安全还可以进一步定义当多个线程访问某个类时不管运行时环境采用何种调度方式或者这些线程将如何交替进行并且在主调代码中不需要任何额外的同步或协同这个类都能表现出正确的行为那么就称这个类是线程安全的。生产者KafkaProducer的实现是线程安全的。KafkaProducer就是一个不可变类。线程安全的可以在多个线程中共享单个KafkaProducer实例所有字段用private final修饰且不提供任何修改方法这种方式可以确保多线程安全。如何节约资源的多线程使用KafkaProducer实例package com.msb.concurrent; ​ import com.msb.selfserial.User; import org.apache.kafka.clients.producer.Callback; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.serialization.StringSerializer; ​ import java.util.Properties; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; ​ /** * 类说明多线程下使用生产者 */ public class KafkaConProducer { ​ //发送消息的个数 private static final int MSG_SIZE 1000; //负责发送消息的线程池 private static ExecutorService executorService Executors.newFixedThreadPool( Runtime.getRuntime().availableProcessors()); private static CountDownLatch countDownLatch new CountDownLatch(MSG_SIZE); ​ private static User makeUser(int id){ User user new User(id); String userName msb_id; user.setName(userName); return user; } ​ /*发送消息的任务*/ private static class ProduceWorker implements Runnable{ ​ private ProducerRecordString,String record; private KafkaProducerString,String producer; ​ public ProduceWorker(ProducerRecordString, String record, KafkaProducerString, String producer) { this.record record; this.producer producer; } ​ public void run() { final String id Thread.currentThread().getId() -System.identityHashCode(producer); try { producer.send(record, new Callback() { public void onCompletion(RecordMetadata metadata, Exception exception) { if(null!exception){ exception.printStackTrace(); } if(null!metadata){ System.out.println(id| String.format(偏移量%s,分区%s, metadata.offset(), metadata.partition())); } } }); System.out.println(id:数据[record]已发送。); countDownLatch.countDown(); } catch (Exception e) { e.printStackTrace(); } } } ​ public static void main(String[] args) { // 设置属性 Properties properties new Properties(); // 指定连接的kafka服务器的地址 properties.put(bootstrap.servers,127.0.0.1:9092); // 设置String的序列化 properties.put(key.serializer, StringSerializer.class); properties.put(value.serializer, StringSerializer.class); // 构建kafka生产者对象 KafkaProducerString,String producer new KafkaProducerString, String(properties); try { for(int i0;iMSG_SIZE;i){ User user makeUser(i); ProducerRecordString,String record new ProducerRecordString,String(concurrent-test,null, System.currentTimeMillis(), user.getId(), user.toString()); executorService.submit(new ProduceWorker(record,producer)); } countDownLatch.await(); } catch (Exception e) { e.printStackTrace(); } finally { producer.close(); executorService.shutdown(); } } ​ ​ ​ ​ } ​消费者KafkaConsumer的实现不是线程安全的实现消费者多线程最常见的方式线程封闭——即为每个线程实例化一个 KafkaConsumer对象package com.msb.concurrent; ​ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.common.serialization.StringDeserializer; ​ import java.time.Duration; import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Properties; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; ​ /** * 类说明多线程下正确的使用消费者需要记住一个线程一个消费者 */ public class KafkaConConsumer { ​ public static final int CONCURRENT_PARTITIONS_COUNT 2; ​ private static ExecutorService executorService Executors.newFixedThreadPool(CONCURRENT_PARTITIONS_COUNT); ​ private static class ConsumerWorker implements Runnable{ ​ private KafkaConsumerString,String consumer; ​ public ConsumerWorker(MapString, Object config, String topic) { Properties properties new Properties(); properties.putAll(config); this.consumer new KafkaConsumerString, String(properties); consumer.subscribe(Collections.singletonList(topic)); } ​ public void run() { final String ThreadName Thread.currentThread().getName(); try { while(true){ ConsumerRecordsString, String records consumer.poll(Duration.ofSeconds(1)); for(ConsumerRecordString, String record:records){ System.out.println(ThreadName|String.format( 主题%s分区%d偏移量%d key%svalue%s, record.topic(),record.partition(), record.offset(),record.key(),record.value())); //do our work } } } finally { consumer.close(); } } } ​ public static void main(String[] args) { ​ /*消费配置的实例*/ MapString,Object properties new HashMapString, Object(); properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,127.0.0.1:9092); properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class); properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class); properties.put(ConsumerConfig.GROUP_ID_CONFIG,c_test); properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,earliest); ​ for(int i 0; iCONCURRENT_PARTITIONS_COUNT; i){ //一个线程一个消费者 executorService.submit(new ConsumerWorker(properties, concurrent-test)); } } ​ ​ ​ ​ } ​群组协调消费者要加入群组时会向群组协调器发送一个JoinGroup请求第一个加入群主的消费者成为群主群主会获得群组的成员列表并负责给每一个消费者分配分区。分配完毕后群主把分配情况发送给群组协调器协调器再把这些信息发送给所有的消费者每个消费者只能看到自己的分配信息只有群主知道群组里所有消费者的分配信息。群组协调的工作会在消费者发生变化(新加入或者掉线)主题中分区发生了变化增加时发生。组协调器组协调器是Kafka服务端自身维护的。组协调器(GroupCoordinator)可以理解为各个消费者协调器的一个中央处理器, 每个消费者的所有交互都是和组协调器(GroupCoordinator)进行的。选举Leader消费者客户端处理申请加入组的客户端再平衡后同步新的分配方案维护与客户端的心跳检测管理消费者已消费偏移量,并存储至__consumer_offset中kafka上的组协调器(GroupCoordinator)协调器有很多有多少个__consumer_offset分区, 那么就有多少个组协调器(GroupCoordinator)默认情况下,__consumer_offset有50个分区, 每个消费组都会对应其中的一个分区对应的逻辑为 hash(group.id)%分区数。消费者协调器每个客户端消费者的客户端都会有一个消费者协调器,他的主要作用就是向组协调器发起请求做交互, 以及处理回调逻辑向组协调器发起入组请求向组协调器发起同步组请求(如果是Leader客户端,则还会计算分配策略数据放到入参传入)发起离组请求保持跟组协调器的心跳线程向组协调器发送提交已消费偏移量的请求消费者加入分组的流程1、客户端启动的时候, 或者重连的时候会发起JoinGroup的请求来申请加入的组中。2、当前客户端都已经完成JoinGroup之后, 客户端会收到JoinGroup的回调, 然后客户端会再次向组协调器发起SyncGroup的请求来获取新的分配方案3、当消费者客户端关机/异常 时, 会触发离组LeaveGroup请求。当然有主动的消费者协调器发起离组请求也有组协调器一直会有针对每个客户端的心跳检测, 如果监测失败,则就会将这个客户端踢出Group。4、客户端加入组内后, 会一直保持一个心跳线程,来保持跟组协调器的一个感知。并且组协调器会针对每个加入组的客户端做一个心跳监测如果监测到过期, 则会将其踢出组内并再平衡。消费者消费的offset的存储consumer_offsets topic并且默认提供了kafka_consumer_groups.sh脚本供用户查看consumer信息。consumer_offsets 是 kafka 自行创建的和普通的 topic 相同。它存在的目的之一就是保存 consumer 提交的位移。kafka-consumer-groups.bat --bootstrap-server :9092 --group c_test --describe那么如何使用 kafka 提供的脚本查询某消费者组的元数据信息呢Math.abs(groupID.hashCode()) % numPartitions__consumer_offsets 的每条消息格式大致如图所示可以想象成一个 KV 格式的消息key 就是一个三元组group.idtopic分区号而 value 就是 offset 的值分区再均衡当消费者群组里的消费者发生变化或者主题里的分区发生了变化都会导致再均衡现象的发生。从前面的知识中我们知道Kafka中存在着消费者对分区所有权的关系这样无论是消费者变化比如增加了消费者新消费者会读取原本由其他消费者读取的分区消费者减少原本由它负责的分区要由其他消费者来读取增加了分区哪个消费者来读取这个新增的分区这些行为都会导致分区所有权的变化这种变化就被称为再均衡。再均衡对Kafka很重要这是消费者群组带来高可用性和伸缩性的关键所在。不过一般情况下尽量减少再均衡因为再均衡期间消费者是无法读取消息的会造成整个群组一小段时间的不可用。消费者通过向称为群组协调器的broker不同的群组有不同的协调器发送心跳来维持它和群组的从属关系以及对分区的所有权关系。如果消费者长时间不发送心跳群组协调器认为它已经死亡就会触发一次再均衡。心跳由单独的线程负责相关的控制参数为max.poll.interval.ms。消费者提交偏移量导致的问题当我们调用poll方法的时候broker返回的是生产者写入Kafka但是还没有被消费者读取过的记录消费者可以使用Kafka来追踪消息在分区里的位置我们称之为偏移量。消费者更新自己读取到哪个消息的操作我们称之为提交。消费者是如何提交偏移量的呢消费者会往一个叫做_consumer_offset的特殊主题发送一个消息里面会包括每个分区的偏移量。发生了再均衡之后消费者可能会被分配新的分区为了能够继续工作消费者者需要读取每个分区最后一次提交的偏移量然后从指定的地方继续做处理。分区再均衡的例子某软件公司有一个项目有两块的工作有两个码农一个小王、一个小李一个负责一块分区消费干得好好的。突然一天小王桌子一拍不干了老子中了5百万了不跟你们玩了立马收拾完电脑就走了。这个时候小李就必须承担两块工作这个时候就是发生了分区再均衡。过了几天你入职一个萝卜一个坑你就入坑了你承担了原来小王的工作。这个时候又会发生了分区再均衡。1如果提交的偏移量小于消费者实际处理的最后一个消息的偏移量处于两个偏移量之间的消息会被重复处理2如果提交的偏移量大于客户端处理的最后一个消息的偏移量,那么处于两个偏移量之间的消息将会丢失再均衡监听器实战我们创建一个分区数是3的主题rebalancekafka-topics.bat --bootstrap-server localhost:9092 --create --topic rebalance --replication-factor 1 --partitions 3在为消费者分配新分区或移除旧分区时,可以通过消费者API执行一些应用程序代码在调用 subscribe()方法时传进去一个 ConsumerRebalancelistener实例就可以了。ConsumerRebalancelistener有两个需要实现的方法。public void onPartitionsRevoked( Collection TopicPartition partitions)方法会在再均衡开始之前和消费者停止读取消息之后被调用。如果在这里提交偏移量下一个接管分区的消费者就知道该从哪里开始读取了public void onPartitionsAssigned( Collection TopicPartition partitions)方法会在重新分配分区之后和消费者开始读取消息之前被调用。具体使用我们先创建一个3分区的主题然后实验一下在再均衡开始之前会触发onPartitionsRevoked方法在再均衡开始之后会触发onPartitionsAssigned方法

相关新闻

WRFDA资料同化实践技术应用

WRFDA资料同化实践技术应用

数值预报已经成为提升预报质量的重要手段,而模式初值质量是决定数值预报质量的重要环节。资料同化作为提高模式初值质量的有效方法,成为当前气象、海洋和大气环境和水文等诸多领域科研、业务预报中的关键科学方法。资料同化新方法的快速发展,…

2026/8/4 3:57:29 阅读更多 →
架空输电线路优质完整报告123(设计源文件+万字报告+讲解)(支持资料、图片参考_相关定制)_文章底部可以扫码

架空输电线路优质完整报告123(设计源文件+万字报告+讲解)(支持资料、图片参考_相关定制)_文章底部可以扫码

架空输电线路优质完整报告123(设计源文件万字报告讲解)(支持资料、图片参考_相关定制)_文章底部可以扫码 45页,15000字,含MATLAB代码

2026/8/4 3:56:29 阅读更多 →
Sunshine游戏串流完整指南:5步搭建你的私人游戏云平台

Sunshine游戏串流完整指南:5步搭建你的私人游戏云平台

Sunshine游戏串流完整指南:5步搭建你的私人游戏云平台 【免费下载链接】Sunshine Self-hosted game stream host for Moonlight. 项目地址: https://gitcode.com/GitHub_Trending/su/Sunshine Sunshine是一款开源的自托管游戏串流服务器,专为Moon…

2026/8/4 3:56:29 阅读更多 →

最新新闻

Python命令行参数解析:从sys.argv到argparse与click实战指南

Python命令行参数解析:从sys.argv到argparse与click实战指南

1. 从命令行到脚本:为什么参数传递是Python开发的必修课如果你写过Python脚本,尤其是那些需要处理不同输入、配置不同运行模式的脚本,那么你一定遇到过这个问题:如何让脚本“听话”地接收外部指令?是每次打开代码文件修…

2026/8/4 4:40:51 阅读更多 →
坐标系转换核心技术:从原理到多传感器融合实战

坐标系转换核心技术:从原理到多传感器融合实战

1. 坐标系转换:从概念到实战的全面拆解在任何一个涉及空间数据处理的领域,无论是游戏开发、机器人导航、地理信息系统(GIS),还是计算机视觉和三维建模,你都无法绕开一个核心问题:坐标系转换。这…

2026/8/4 4:40:51 阅读更多 →
AI Agent概念辨析与落地实践:从技术维度到应用挑战

AI Agent概念辨析与落地实践:从技术维度到应用挑战

1. 从“智能体”到“智能代理”:一个概念的混乱史 最近和几个不同领域的朋友聊天,发现一个挺有意思的现象:当大家提到“AI Agent”这个词时,脑子里想的完全不是一回事儿。搞自动驾驶的哥们儿,觉得Agent就是那套能感知、…

2026/8/4 4:40:51 阅读更多 →
分布滞后模型:从原理到实战,解析时间序列中的动态影响

分布滞后模型:从原理到实战,解析时间序列中的动态影响

1. 从“昨天”到“明天”:理解滞后效应的现实场景在商业分析、经济研究和政策评估中,我们常常会遇到一个看似简单却极其棘手的问题:一个事件的影响,往往不是立竿见影的。比如,公司今天投入一笔营销费用,销售…

2026/8/4 4:40:51 阅读更多 →
【Azure APIM】通过 API Management 公开现有 MCP Server 的试验 (一)

【Azure APIM】通过 API Management 公开现有 MCP Server 的试验 (一)

随着 GitHub Copilot、Claude、ChatGPT 等 AI 客户端逐渐支持 Model Context Protocol(MCP),企业需要以统一、可控的方式向这些客户端提供内部工具。 对于已经存在的远程 MCP Server,如果直接让客户端访问,通常难以统一…

2026/8/4 4:40:51 阅读更多 →
Cocos Creator内存泄漏排查实战:从工具使用到典型场景解析

Cocos Creator内存泄漏排查实战:从工具使用到典型场景解析

1. 项目概述:为什么Cocos内存泄漏是开发者的“心腹大患”?干了这么多年游戏开发,尤其是用Cocos Creator,最让人头疼的往往不是炫酷的特效做不出来,而是游戏跑着跑着就卡了、闪退了,或者玩家手机发烫、电量狂…

2026/8/4 4:39:51 阅读更多 →

日新闻

AI Agent白手起家26: 使用标准事件驱动大模型实践

AI Agent白手起家26: 使用标准事件驱动大模型实践

纲要 练习目标:掌握大模型标准事件的调用回顾 LangChain 中的核心标准事件 invokestreambatchastream_eventswith_structured_output 环境准备实战代码:多种事件调用对比 同步调用与流式输出批量处理异步事件流监听结构化输出 运行说明与预期结果总结与扩…

2026/8/4 0:00:40 阅读更多 →
dealsea是什么?跨境卖家必知的美国deal站入门指南

dealsea是什么?跨境卖家必知的美国deal站入门指南

说实话,第一次听说美国这个老牌折扣网站的跨境卖家,十个有八个会问同一个问题:这个平台到底是干嘛的?我见过一个做家居出口的朋友,他在亚马逊上月销二十万美金,却从来没用过它。我给他看了首页——一屏一屏…

2026/8/4 0:01:40 阅读更多 →
清华大学重磅EST:植物自导电闪蒸焦耳热600°C/2600°C两步法!稀土超积累植物秒级转化为CeO₂-石墨烯电催化剂!

清华大学重磅EST:植物自导电闪蒸焦耳热600°C/2600°C两步法!稀土超积累植物秒级转化为CeO₂-石墨烯电催化剂!

通讯作者:邓兵、刘建国通讯单位:清华大学DOI:https://doi.org/10.1021/acs.est.6c00603研究背景稀土元素(REEs)是清洁能源技术与电子器件不可或缺的核心原料,然而传统提取方式依赖能耗高、排放大的采矿与强…

2026/8/4 0:01:40 阅读更多 →

周新闻

最大流算法详解:从水管网络到Ford-Fulkerson与Dinic实战

最大流算法详解:从水管网络到Ford-Fulkerson与Dinic实战

1. 从水管网络到最大流:一个核心问题的诞生想象一下,你是一个城市供水系统的总工程师。你的城市有多个水源(水库),需要通过一个复杂的地下管道网络,将水输送到各个居民区。每条管道都有其最大通水能力&…

2026/8/3 4:58:13 阅读更多 →
基于Springboot的企业门户网站(源码+LW+调试文档+讲解)

基于Springboot的企业门户网站(源码+LW+调试文档+讲解)

温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台…

2026/8/3 1:53:31 阅读更多 →
MATLAB xcorr函数详解:从互相关原理到四大实战应用

MATLAB xcorr函数详解:从互相关原理到四大实战应用

1. 从一次信号“找茬”说起:为什么我们需要互相关几年前,我在处理一组声学传感器数据时遇到了一个棘手的问题。我有两个麦克风记录了一段相同的音频信号,理论上它们接收到的声音波形应该非常相似,只是由于麦克风位置不同&#xff…

2026/8/3 4:36:35 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/3 13:07:03 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/3 5:19:38 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/3 8:27:36 阅读更多 →