首页 文章 精选 留言 我的

精选列表

搜索[心智成长],共7680篇文章
优秀的个人博客,低调大师

物联网产业7年 变局中不断成长

物联网产业的媒体联动原素创办者李晓妍自2011年开始从事物联网媒体以来,目睹了中国物联网产业在最初这几年的跌宕起伏。此文就是她对物联网产业7年变局的见证。 翻阅物联网的相关资料,基本都显示“物联网”这一概念,是1995年比尔·盖茨在其《未来之路》一书中提出的。但是,物联网的“产业化概念”却起源于中国,这一点,想必全世界都会认同。 2009年8月,中国前总理温家宝在一次无锡视察中,听取了无锡物联网产业研究院负责人对“物联网产业化”的一些建议后,表示了认同,并提出了三条具体实施建议。此后,让“物联网”从学术概念走向了产业概念。无锡,也成为当之无愧的物联网产业发祥地。 截至2016年8月,物联网产业化概念整整走过了7年的历程。婚姻中有“7年之痒”一说,可以理解为7年是一个磨合期,跨过了这道坎,可能就离圆满结局更近一步;跨不过,或许就要一拍两散了。有时候感觉,这竟是对中国物联网产业发展历程的绝佳映射。因为中国物联网产业过去7年的发展,用“迂回曲折”已不足以形容。 物联网给无锡带来的喜与悲 首先,来说说物联网产业的发祥地——无锡。可以说,很多早期的物联网从业者,对无锡是很有感情的,也包括我。因为无锡是大部分从业者认知物联网的开始。但是应该说,无锡成就了物联网,物联网也成就了无锡;同时,无锡毁了物联网,物联网也毁了无锡。这个怎么解释呢?在中国,只要知道物联网的人,都知道它的发祥地是无锡,但接着就会说:“无锡代表的是传统物联网”,尤其是2014年以后参与到物联网产业的创业者们,更是表示“一提物联网,大家首先想到的就是智能抄表、RFID、追踪溯源等,所以我们更愿意提万物互联”。这可能就是作为开山鼻祖必须背的“黑锅”吧! 为什么说是黑锅呢?因为作为中国物联网产业的发祥地,无锡承担了很大一部分的“产业化试验”工作和责任。但是,当时很多产业化的条件都不成熟。 第一,物联网的技术路线还不清晰。是三层架构,还是四层架构;是类似于美国的M2M,仅包含传感层,还是包含了公网及更大范围的信息化系统集成等依然在争论中。 第二,没有清晰切入点,所以干脆“囫囵吞枣”。也就是说,大部分从业者都知道物联网带来的终极体验是:所有人和物都能够通过网络无障碍互联互通,却不知道该从哪儿下手。于是,就把想象中的这张终极大网,切分成无数个小系统,进行试验,名为物联网产业示范试点工程。中国最初的三个物联网示范工程——感知市民中心、智能机场安检等,便是由无锡于2010年1月推行实施的。此后,这一举措还被中央采纳,在全国推广,并且一直延续至今年。 第三,物联网的很多支撑技术和产品尚不成熟。比如传感器技术,当时智能传感器的概念在中国刚刚开始,掌握MEMS工艺的中国本土传感器企业寥寥无几。再比如短距离无线传输技术,RFID、NFC、ZigBee、Bluetooth、Wifi、PLC等都有参与其中,并且都曾在不同时段“艳压”过其它技术,但各自都有很大局限。像RFID的成本与应用场景、ZigBee的标准不统一、Wifi的高功耗,并且当时Wifi刚刚在中国解禁、PLC的干扰性大等问题,无一不对物联网产业的发展造成阻碍。再比如IT技术的云计算化升级也处于初期。自身都还处于混乱中,更别说支持物联网的发展了。还有当时大数据的概念还处于学术阶段,因为其技术基础Hadoop的发展才刚刚开始。 第四,大部分从业者的急于求成。2008金融海啸过后,中国企业太想快速找到下一个经济增长点。因而,当时房地产行业疯狂发展,中国的房价开涨都始于2008年。在那样的环境影响下,很多物联网从业者都希望物联网产业也能够像房地产一样,可以让其“一夜暴富”。这从一定程度上让物联网产业走上了一条不健康的发展路线。 而这些条件不成熟造成的后果就是:1)物联网产业没有像大家期待的那样,瞬间爆发;2)无锡没有像大家期许的那样,瞬间爆发;3)物联网企业,没有像自己想象的那样,瞬间爆发。而无锡这座——曾经因物联网而风光无限的城市,承担了这个后果的部分原罪。并且,它当时的付出,比如通过“530计划”为物联网创业者和企业提供的各种扶持,似乎都成了笑话。 不过,无锡并没有因此就放弃物联网产业。2013年以后,无锡看似不再大造声势了,但是却依然在暗暗发力,更或者是在向着更加扎实的方向探索。尤其是2016年,在中国被称为真正的“物联网产业发展之年”,无锡不仅提出了打造“物联网小镇”的行动方案,而且将第七届物联网博览会升级为“世界物联网博览会”,再次成为物联网产业各界集聚的核心。 2009-2013年:政府主导下的实验室阶段 其次,再看看这7年间,中国物联网产业的进阶历程。我一直比较倾向于将其分为3个阶段。第一个阶段(2009-2013年),可以称为“政府主导下的实验室”阶段,也就是无锡占据绝对主导地位的阶段。这个阶段又可以分为两个小节。 第一节是2009-2011年,在国务院指导下,中央各部委,及各省市纷纷出台如物联网专项基金扶持、物联网示范试点工程、税收优惠等“推动物联网产业发展”的政策。其中的“物联网应用试点示范工程(项目)”,应该说为物联网的产业化从设想到落地提供了试验田。据不完全统计,在2010年至2016年间,中央各部委及地方省市规划及打造的物联网示范试点工程逾千个,涉及项目过万个。当然,这着实养活了一批物联网企业,尤其在早期。记得2012年走访一家物联网企业的时候,该公司有一个叫做“企业发展部”的部门,该部门的主要工作是研究国家政策,申报各种扶持基金和项目,当时这是该公司的主要收入来源。 第二节是2012-2013年,中国政府出台了800亿物联网产业“造市”计划——打造智慧城市。中国住建部分别于2012年底和2013年下半年公布了两批智慧城市试点城市,共193个。每个试点城市都分批规划建设试点项目。瞬间又有过万个“智慧”项目启动。当时的物联网企业就是靠着这样的支撑一步步走下去的。 2014-2016年:局部市场主导阶段 第二个阶段(2014-2016年),局部进入市场主导阶段。但是,这个阶段开始的方式却颇为耐人寻味,始于2014年1月的Google收购Nest。这个事件成为了中国物联网产业发展最明显的分水岭。同时,也是中国物联网力量与国际物联网力量的首次会合。 当然,至今还有一部分人认为,物联网概念在其它国家并不盛行。但是就我看到的状况是,全球在信息科技领域占据主导地位的国家,在物联网产业均有布局,其中走在前列的应该是中国和美国,只是各自的提法不同,开始的切入点也不同。比如美国企业的切入方式更偏向于自上而下,即从支撑技术、产品开始;而中国企业一开始就选择自下而上,从应用开始。 这背后的原因之一可能要归结于,大部分中国人对物联网的理解是:它并非一向新技术,而是将既有技术的整合应用。但是总体来说,大家都从不同角度进行着对物联网的探索。Google收购Nest之后,或许让更多人看清了实现物联网的第一步——硬件智能化。 在此之前,中国参与物联网产业的都是大型企业,如电信运营商、从事B2B业务的大型信息化集成商等,即便是创业企业也是挂靠在政府下面的三产或者四产等等。就如前面所说,当时对物联网应用的理解,就是一个完整系统,大家都争着要做“物联网服务运营商”,并且当时所有能赚钱的物联网项目都是政府扶持的,所以资本不够雄厚、没有资源背景的中小企业,或者创业者根本没机会,而从事B2C业务的互联网公司们又看不上。 但是,Google收购Nest之后,中国的互联网巨头们好像瞬间觉醒了,开始频繁出现于物联网的各种活动和论坛上。同时,原来互联网圈子的投资机构也开始活跃于物联网圈子了。这给草根创业者们带来了前所未有的机遇。 因此,2014年,在物联网产业中一直探讨的“颠覆”出现了第一个迹象:资源移位。原本没有机会的,仅掌握着20%资本和资源,却占据80%人群份额的“草根”们,成为了物联网产业的主力军;在资本的煽动下,各种千奇百怪、异想天开的智能终端创意充斥于市场上;传统制造业,甚至传统IT产业都看不懂市场上在发生着什么,瞬间充满了危机感。这个时候,物联网看上去,终于配的起发起者对它最初的设想了。毕竟,资源分配都不改变,怎么能谈的上颠覆?! 也是在这一年,深圳成为了物联网产业的焦点,因为80%的智能硬件创业者都集聚在深圳。深圳竟然瞬间成为了中国科技城市的代表。记得2015年一本名为《迪拜视野》的杂志,列举了全球最具代表性的8个科技城市,其中,中国上榜的便是深圳。并且据说2016年拉斯维加斯的CES展上,四分之一的企业都来自中国深圳。 但是这个阶段,中国物联网产业取得的突破,都集中在人们生活相关的领域,比如今年中国新上市的白色家电100%实现了智能化;黑色家电40%实现了智能化;小型家电正在逐步走向智能化;小型化、便携式医疗设备大规模出现;智能家居产业火热,虽然仅停留在产业界。而在市政、民生、工业、农业等,原本买方就是政府,或企业的领域,突破依然有限。数月前在走访一家农业物联网企业时,了解到2016年开始,中国政府取缔了对农业物联网项目的资金扶持,农业物联网的推广便进入了举步维艰的地步。 2016年后:物联网迎来第二个春天 第三个阶段(2016年7月以后),整个物联网产业进入良性发展阶段。当然,说起发展,每个阶段其实都在发展。但是以前的发展应该说太过跌宕,太过迂回,太过不健康,所以不能算真正的发展。而2016年,应该说是各种条件都相对成熟后的发展,所以谓“良性发展”。 2016年7月,发生了一件震惊科技圈的大事件:软银322亿美元收购了全球领先的手机芯片设计厂商Arm。软银选择在这个时候收购Arm的原因,并非因为其是手机芯片领域的龙头,而是因为其打造了一款物联网操作系统——ArmEmbeddedOS。也就是说,软银意图的是整个物联网时代。 这个事件与其说是引发物联网产业第三次进阶的导火索,不如说是物联网产业过去7年探索的结果。因为,2016年实现物联网的很多条件和变量几乎都完备了。 第一,物联网的基石——云计算,技术基本达到成熟,并已经形成鲜明产业格局; 第二,物联网的应用支撑——大数据,底层技术架构已经相对稳定,Hadoop的地位已经确立,在某些领域如金融、环保、电力等,已经有成熟应用案例; 第三,物联网传输层的核心技术——NB-IOT和5G可以商用,对整个物联网产业的发展来讲,是一个极大地跨越; 第四,制造业的觉醒。这应该说对物联网产业的发展极其重要。因为物联网的核心是“物”,如果“物”的生产者和制造者不接受新技术,不接受变革,那么将拖慢整个产业的发展速度。 更重要地是,这一年,中国的物联网产业在工业与民用领域开始齐头并进,而不是像2013年前只重工业应用,也不是像2014、2015年间,民用独领风骚。这背后,应该说是政府力量和市场力量的共同推动。比如最近两年,中国政府针对工业4.0、智能制作推出了一系列扶持政策,以及示范试点项目。因为,工业不改革,物联网的应用便不会深入。就像互联网时代,无论说它怎么颠覆了传统产业,最终也只是改变了信息传播通道,以及消费领域的交易渠道。而工业占据着中国80%的经济份额,所以工业领域的变革,将带来的价值远远大于消费领域。总之,这个阶段,在中国,政府与市场开始合力推动物联网产业的发展。 不过,由于中国的经济特色,核心工业一直由政府主导,所以工业领域的变革开始后,或许将再次刷新中国科技城市的排名。当然,这一次,不仅无锡和深圳,很多其它城市如上海、杭州、南京、北京、成都等都已对物联网“虎视眈眈”,所以最终谁能胜出,很难评判。不过,从目前的状况来看,无锡的决心与力度似乎更大一些。 本文转自d1net(转载)

优秀的个人博客,低调大师

🔥 神级成长性!Solon v2.9 发布

Solon 框架! Java “纯血国产”应用开发框架。开放原子开源基金会,孵化项目。从零开始构建(非 java-ee 架构),有灵活的接口规范与开放生态。 追求: 更快、更小、更简单 提倡: 克制、简洁、高效、开放、生态 有什么特点? 特点 描述 更高的计算性价比 并发高 300%;内存省 50% 更快的开发效率 代码少;入门快;调试重启快 10 倍 更好的生产与部署体验 打包小 90% 更大的兼容范围 非 java-ee 架构;同时支持 java8 ~ java22,graalvm native image 入门探索视频(用户录制): 最近更新了什么? 新增 solon.boot.vertx 插件 新增 solon.cloud.gateway 插件 新增 solon.rx 插件 添加 solon.data 配置节solon.dataSources(用于自动构建数据源),支持 ENC 加密符 添加 solon.docs 配置节solon.docs(用于自动构建文档摘要) 添加 solon.view.prefix 配置项支持 "file:" 前缀(支持体外目录) 添加 solon.scheduling.simple SimpleScheduler::isStarted 方法 添加 solon@Condition(onBean, onBeanName)条件属性 添加 solon.validation ValidUtils 工具类 添加 solon LifecycleBean:postStart 方法 添加 solon MethodInterceptor 接口,替代 Interceptor(旧接口保留) 添加 solon.net.httputils 扩展机制,并与 solon.cloud 自动整合 添加 solon.net.httputils HttpResponse::headerNames 方法 添加 solon.cloud CloudDiscoveryService:findServices 方法 添加 solonsolon.plugin.exclude应用属性配置 添加 solonsolon.app.enabled应用属性配置(Solon.cfg().appEnabled()可获取) 添加 solon${.url}应用属性配置本级引用 添加 solon--cfg启动参数支持(便于内嵌场景开发) 添加 托管类构造参数注入支持(对 kotlin 更友好) 调整 solon.cloud.httputils 标为弃用,由 solon.net.httputils 替代 调整 smarthttp,jetty,undertow 的非标准方法的 FormUrlencoded 预处理时机 调整 solon.auth maven 包更名为 solon.security.auth (原 maven 包保留) 调整 solon.validation maven 包更名为 solon.security.validation (原 maven 包保留) 调整 solon.vault maven 包更名为 solon.security.vault (原 maven 包保留) 优化 AppContext::beanMake 保持与 beanSacn 相同的类处理 优化 solon.serialization.jackson 兼容 @JsonFormat 注解时间格式和时间格式配置并存 优化 solon Context::body 的兼容性,避免不可读情况 优化 solon 调试模式与 gradle 的兼容性 优化 solon.boot FormUrlencodedUtils 预处理把 post 排外 优化 solon.web.rx 允许多次渲染输出 优化 kafka-solon-cloud-plugin 添加 username, password 简化配置支持(简化有账号的连接体验) 优化 solon.boot 413 状态处理 优化 solon.boot.smarthttp 适配的 maxRequestSize 设置(取 fileSize 和 bodySize 的大值) 优化 solon AppContext 注册和查找时以 rawClz 为主(避免以接口注册时,实例类型查不到) 优化 solon.mvc kotlin data class 带默认值的注入支持(表单模式下) 优化 solon PathAnalyzer 添加 addStarts 参数选择,支持域名匹配 优化 solon LifecycleBean 和 Lifecycle 设计 修复 solon.view.thymeleaf 模板不存在时没有输出 500 的问题 修复 solon.serialization.jackson 泛型注入失效的问题 修复 solon.boot.smarthttp 适配在 chunked 下不能读取 body string 的问题 修复 solon-openapi2-knife4j 没有配置时不能启动的问题(默认改为不启用) wood 升为 1.3.0 snack3 升为 3.2.109 socket.d 升为 2.5.11 zookeeper 升为 3.9.2 dromara-plugins 升为 0.1.2 kafka_2.13 升为 3.8.0 beetlsql 升为 3.30.10-RELEASE beetl 升为 3.17.0.RELEASE mybatis 升为 3.5.16 mybatis-flex 升为 1.9.6 sqltoy 升为 5.6.20 dbvisitor 升为 5.4.3 bean-searcher 升为 4.3.0 liteflow 升为 2.12.2 aws.s3 升为 1.12.769 powerjob 升为 5.1.0 netty 升为 4.1.112.Final reactor-core 升为 3.6.9 reactor-netty-http 升为 1.1.22 vertx 升为 4.5.9 undertow 升为 2.2.34.Final jetty 升为 9.4.55.v20240627 smarthttp 升为 1.5.9 项目仓库地址? gitee:https://gitee.com/opensolon/solon github:https://github.com/opensolon/solon 官网? https://solon.noear.org

优秀的个人博客,低调大师

Kafka成长3:Producer 元数据拉取源码原理(上)

上一节我们分析了Producer的核心组件,我们得到了一张关键的组件图。你还记得么? 简单概括下上面的图就是: 创建了Metadata组件,内部通过Cluster维护元数据 初始化了发送消息的内存缓冲器RecordAccumulator 创建了NetworkClient,内部最重要的是创建了NIO的Selector组件 启动了一个Sender线程,Sender引用了上面的所有组件,开始执行run方法。 图的最下方可以看到,上一节截止到了run方法的执行,这一节我们首先会看看run方法核心脉络做了什么。接着分析下Producer第一个核心流程:元数据拉取的源码原理。 让我们开始吧! Sender的run方法在做什么? 这一节我们就继续分析下,sender的run方法开始执行会做什么。 public void run() { log.debug("Starting Kafka producer I/O thread."); // main loop, runs until close is called while (running) { try { run(time.milliseconds()); } catch (Exception e) { log.error("Uncaught error in kafka producer I/O thread: ", e); } } log.debug("Beginning shutdown of Kafka producer I/O thread, sending remaining records."); // okay we stopped accepting requests but there may still be // requests in the accumulator or waiting for acknowledgment, // wait until these are completed. while (!forceClose && (this.accumulator.hasUnsent() || this.client.inFlightRequestCount() > 0)) { try { run(time.milliseconds()); } catch (Exception e) { log.error("Uncaught error in kafka producer I/O thread: ", e); } } if (forceClose) { // We need to fail all the incomplete batches and wake up the threads waiting on // the futures. this.accumulator.abortIncompleteBatches(); } try { this.client.close(); } catch (Exception e) { log.error("Failed to close network client", e); } log.debug("Shutdown of Kafka producer I/O thread has completed."); } 这个run方法的核心脉络很简单。主要就是2个while循环+线程的close,而2个while循环,他们都调用了run(long time)的这个方法。 通过注释你可以看到,第二个while是处理特殊情况的,当第一个while退出后,还有未发送的请求,需要第二个while循环处理完成,才会关闭线程。 整体脉络如下图所示: 接着其实就该看下run方法主要在干什么了? /** * Run a single iteration of sending * * @param now * The current POSIX time in milliseconds */ void run(long now) { Cluster cluster = metadata.fetch(); // get the list of partitions with data ready to send RecordAccumulator.ReadyCheckResult result = this.accumulator.ready(cluster, now); // if there are any partitions whose leaders are not known yet, force metadata update if (result.unknownLeadersExist) this.metadata.requestUpdate(); // remove any nodes we aren't ready to send to Iterator<Node> iter = result.readyNodes.iterator(); long notReadyTimeout = Long.MAX_VALUE; while (iter.hasNext()) { Node node = iter.next(); if (!this.client.ready(node, now)) { iter.remove(); notReadyTimeout = Math.min(notReadyTimeout, this.client.connectionDelay(node, now)); } } // create produce requests Map<Integer, List<RecordBatch>> batches = this.accumulator.drain(cluster, result.readyNodes, this.maxRequestSize, now); if (guaranteeMessageOrder) { // Mute all the partitions drained for (List<RecordBatch> batchList : batches.values()) { for (RecordBatch batch : batchList) this.accumulator.mutePartition(batch.topicPartition); } } List<RecordBatch> expiredBatches = this.accumulator.abortExpiredBatches(this.requestTimeout, now); // update sensors for (RecordBatch expiredBatch : expiredBatches) this.sensors.recordErrors(expiredBatch.topicPartition.topic(), expiredBatch.recordCount); sensors.updateProduceRequestMetrics(batches); List<ClientRequest> requests = createProduceRequests(batches, now); // If we have any nodes that are ready to send + have sendable data, poll with 0 timeout so this can immediately // loop and try sending more data. Otherwise, the timeout is determined by nodes that have partitions with data // that isn't yet sendable (e.g. lingering, backing off). Note that this specifically does not include nodes // with sendable data that aren't ready to send since they would cause busy looping. long pollTimeout = Math.min(result.nextReadyCheckDelayMs, notReadyTimeout); if (result.readyNodes.size() > 0) { log.trace("Nodes with data ready to send: {}", result.readyNodes); log.trace("Created {} produce requests: {}", requests.size(), requests); pollTimeout = 0; } for (ClientRequest request : requests) client.send(request, now); // if some partitions are already ready to be sent, the select time would be 0; // otherwise if some partition already has some data accumulated but not ready yet, // the select time will be the time difference between now and its linger expiry time; // otherwise the select time will be the time difference between now and the metadata expiry time; this.client.poll(pollTimeout, now); } 上面的代码,你如果第一次看,你肯定会觉得,这个脉络非常不清晰,不知道重点在哪里。不过还好有些注释,你能大体猜到他在干嘛。 accumulator的ready,networkclient的ready、networkclient的send、networkclient的poll 这些好像是在准备内存区域、准备网络连接的node节点、发送数据、拉取响应结果的意思。 可是如果你猜不到,该怎么办呢? 这时候就可以祭出debug这个杀器了。由于是producer,我们可以在Hellowolrd的这个客户端打断点,一步一步看下。 当你对run方法一步一步打了断点之后你会发现: accumulator的ready,networkclient的ready、networkclient的send 这些的逻辑几乎都没有执行,全部都是初始化空对象,或者方法内部直接return。 直接一路执行到了client.poll方法。如下图所示: 那么,你可以得出一个结论,while第一次循环这个run方法的核心逻辑,其实只有一句话: client.poll(pollTimeout, now) 整体脉络如下所示: 看来接下来,这个NetworkClient的poll方法,就是关键中的关键了: /** * Do actual reads and writes to sockets. * 对套接字进行实际读取和写入 * * @param timeout The maximum amount of time to wait (in ms) for responses if there are none immediately, * must be non-negative. The actual timeout will be the minimum of timeout, request timeout and * metadata timeout * @param now The current time in milliseconds * @return The list of responses received */ @Override public List<ClientResponse> poll(long timeout, long now) { long metadataTimeout = metadataUpdater.maybeUpdate(now); try { this.selector.poll(Utils.min(timeout, metadataTimeout, requestTimeoutMs)); } catch (IOException e) { log.error("Unexpected error during I/O", e); } // process completed actions long updatedNow = this.time.milliseconds(); List<ClientResponse> responses = new ArrayList<>(); handleCompletedSends(responses, updatedNow); handleCompletedReceives(responses, updatedNow); handleDisconnections(responses, updatedNow); handleConnections(); handleTimedOutRequests(responses, updatedNow); // invoke callbacks for (ClientResponse response : responses) { if (response.request().hasCallback()) { try { response.request().callback().onComplete(response); } catch (Exception e) { log.error("Uncaught error in request completion:", e); } } } return responses; } 这个方法的脉络就清晰多了,通过方法名和注释,我们几乎可以猜出他的一些作用主要有: 1)注释说:对套接字进行实际读取和写入 2)metadataUpdater.maybeUpdate(),你还记得NetworkClient的组件DefaultMetadataUpdater么,方法名意思是可能进行元数据更新。这个好像很关键的样子 3)接着执行了Selector的poll方法,这个是NetworkClient的另一个组件Selector,还记得么?它底层封装了原生的NIO Selector。这个方法应该也比较关键。 4)后续对response执行了一系列的方法,从名字上看, handleCompletedSends 处理完成发送的请求、handleCompletedReceives处理完成接受的请求、handleDisconnections处理断开连接的请求、handleConnections处理连接成功的请求、处理超时的请求handleTimedOutRequests。根据不同情况有不同的处理。 5)最后还有一个response的相关的回调处理,如果注册了回调函数,会执行下。这个应该不是很关键的逻辑 也就是简单的说就是NetworkClient执行poll方法,主要通过selector处理请求的读取和写入,对响应结果做不同的处理而已。 如下图所示: 到这里其实我们基本摸清出了run方法主要在做的一件事情了,由于是第一次循环,之前的accumulator的ready,networkclient的ready、networkclient的send 什么都没做,第一次while循环run方法核心执行的是networkclient.poll方法。而poll方法的主要逻辑就是上面图中所示的了。 maybeUpdate可能在在拉取元数据? 刚才我们分析到,poll方法首先执行的是DefaultMetadataUpdater的maybeUpdate方法,它是可能更新的意思。我们来一起看下他的逻辑吧。 public long maybeUpdate(long now) { // should we update our metadata? long timeToNextMetadataUpdate = metadata.timeToNextUpdate(now); long timeToNextReconnectAttempt = Math.max(this.lastNoNodeAvailableMs + metadata.refreshBackoff() - now, 0); long waitForMetadataFetch = this.metadataFetchInProgress ? Integer.MAX_VALUE : 0; // if there is no node available to connect, back off refreshing metadata long metadataTimeout = Math.max(Math.max(timeToNextMetadataUpdate, timeToNextReconnectAttempt), waitForMetadataFetch); if (metadataTimeout == 0) { // Beware that the behavior of this method and the computation of timeouts for poll() are // highly dependent on the behavior of leastLoadedNode. Node node = leastLoadedNode(now); maybeUpdate(now, node); } return metadataTimeout; } /** * The next time to update the cluster info is the maximum of the time the current info will expire and the time the * current info can be updated (i.e. backoff time has elapsed); If an update has been request then the expiry time * is now */ public synchronized long timeToNextUpdate(long nowMs) { long timeToExpire = needUpdate ? 0 : Math.max(this.lastSuccessfulRefreshMs + this.metadataExpireMs - nowMs, 0); long timeToAllowUpdate = this.lastRefreshMs + this.refreshBackoffMs - nowMs; return Math.max(timeToExpire, timeToAllowUpdate); } 原来这里有一个时间的判断,当判断满足才会执行maybeUpdate。 这个时间计算好像比较复杂,但是大体可以看出来,metadataTimeout是根据三个时间综合判断出来的,如果是0才会执行真正的maybeUpdate()。 像这种时候,我们可以直接在metadataTimeout这里打一个断点,看下它的值是如何计算的,比如下图: 你会发现,当第一次执行while循环,执行到poll方法,执行到这个maybeUpdate的时候,决定metadataTimeout的3个值,有两个是0,其中一个是非0,是一个299720的值。最终导致metadataTimeout也是非0,是299720。 也就是说,第一次while循环不会执行maybeUpdate的任何逻辑。 那么接着向下执行 Selector的poll()方法。 /** * Do whatever I/O can be done on each connection without blocking. This includes completing connections, completing * disconnections, initiating new sends, or making progress on in-progress sends or receives. * 在不阻塞的情况下,在每个连接上做任何可以做的 I/O。这包括完成连接完成、断开连接,启动新的发送,或在进行中的发送或接收请求 */ @Override public void poll(long timeout) throws IOException { if (timeout < 0) throw new IllegalArgumentException("timeout should be >= 0"); clear(); if (hasStagedReceives() || !immediatelyConnectedKeys.isEmpty()) timeout = 0; /* check ready keys */ long startSelect = time.nanoseconds(); //这个方法是NIO底层Selector.select(),会阻塞监听 int readyKeys = select(timeout); long endSelect = time.nanoseconds(); currentTimeNanos = endSelect; this.sensors.selectTime.record(endSelect - startSelect, time.milliseconds()); //如果监听到有操作的SelectionKeys,也就是readyKeys>0< 会执行一些操作 if (readyKeys > 0 || !immediatelyConnectedKeys.isEmpty()) { pollSelectionKeys(this.nioSelector.selectedKeys(), false); pollSelectionKeys(immediatelyConnectedKeys, true); } addToCompletedReceives(); long endIo = time.nanoseconds(); this.sensors.ioTime.record(endIo - endSelect, time.milliseconds()); maybeCloseOldestConnection(); } private int select(long ms) throws IOException { if (ms < 0L) throw new IllegalArgumentException("timeout should be >= 0"); if (ms == 0L) return this.nioSelector.selectNow(); else return this.nioSelector.select(ms); } 上面的脉络主要是2步: 1)select(timeout): NIO底层selector.select(),会阻塞监听 2)pollSelectionKeys(): 监听到有操作的SelectionKeys,做了一些操作 也就是说,最终,Sender线程的run方法,第一次while循环执行poll方法,最后什么都没干,会被selector.select()阻塞住。 如下图所示: new KafkaProducer之后 分析完了run方法的执行 ,我们分析的KafkaProducerHelloWorld第一步new KafkaProducer()基本就完成了。 大家经历了一节半的时间,终于分析清楚了KafkaProducer创建的原理。不不知道你对Kafka的Producer是不是有了更深的理解了。 分析了new KafkaProducer()之后呢? 我们继续接着KafkaProducerHelloWorld往下分析,你还记得KafkaProducerHelloWorld的代码么? public class KafkaProducerHelloWorld { public static void main(String[] args) throws Exception { //配置Kafka的一些参数 Properties props = new Properties(); props.put("bootstrap.servers", "mengfanmao.org:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 创建一个Producer实例 KafkaProducer<String, String> producer = new KafkaProducer<>(props); // 封装一条消息 ProducerRecord<String, String> record = new ProducerRecord<>( "test-topic", "test-key", "test-value"); // 同步方式发送消息,会阻塞在这里,直到发送完成 // producer.send(record).get(); // 异步方式发送消息,不阻塞,设置一个监听回调函数即可 producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if(exception == null) { System.out.println("消息发送成功"); } else { exception.printStackTrace(); } } }); Thread.sleep(5 * 1000); // 退出producer producer.close(); } KafkaProducerHelloWorld主要就3步: 1)new KafkaProducer 这个我们已经分析完了,主要分析了配置文件的解析、各个组件是什么、有什么,还有就是刚才分析的run线程第一次循环到底执行了什么。 2) new ProducerRecord 创建待发送的消息 3) producer.send() 发送消息 首先创建待发送的消息: ProducerRecord<String, String> record = new ProducerRecord<>("test-topic", "test-key", "test-value"); public ProducerRecord(String topic, K key, V value) { this(topic, null, null, key, value); } /** * Creates a record with a specified timestamp to be sent to a specified topic and partition * 创建具有指定时间戳的记录以发送到指定主题和分区 * @param topic The topic the record will be appended to * @param partition The partition to which the record should be sent * @param timestamp The timestamp of the record * @param key The key that will be included in the record * @param value The record contents */ public ProducerRecord(String topic, Integer partition, Long timestamp, K key, V value) { if (topic == null) throw new IllegalArgumentException("Topic cannot be null"); if (timestamp != null && timestamp < 0) throw new IllegalArgumentException("Invalid timestamp " + timestamp); this.topic = topic; this.partition = partition; this.key = key; this.value = value; this.timestamp = timestamp; } 我们之前提过,Record表示了一条消息的抽象封装。这个ProducerRecord其实就表示了一条消息。 从构造函数的注释可以看出来**,ProducerRecord可以指定往哪个topic,哪一个分区partition,并且消息可以设置一个时间戳。分区和时间戳默认可以不指定** 其实看这块源码,我们主要得到的信息就是这些了,这些都比较简单。就不画图了。 发送消息时的元数据拉取触发 当Producer和Record都创建好了之后,可以用同步或者异步的方式发送消息。 // 同步方式发送消息,会阻塞在这里,直到发送完成 // producer.send(record).get(); // 异步方式发送消息,不阻塞,设置一个监听回调函数即可 producer.send(record, new Callback() { @Override public void onCompletion(RecordMetadata metadata, Exception exception) { if(exception == null) { System.out.println("消息发送成功"); } else { exception.printStackTrace(); } } }); //同步发送 @Override public Future<RecordMetadata> send(ProducerRecord<K, V> record) { return send(record, null); } //异步发送 public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) { // intercept the record, which can be potentially modified; this method does not throw exceptions ProducerRecord<K, V> interceptedRecord = this.interceptors == null ? record : this.interceptors.onSend(record); return doSend(interceptedRecord, callback); } 同步和异步的整个发送逻辑如下图所示: 从上图你会发现,但是无论同步发送还是异步底层都会调用同一个方法doSend()。区别就是有没有callBack回调函数而已,他们还都在调用前注册一些拦截器,这里我们抓大放小下,我们重点还是关注doSend方法。 doSend方法如下: /** * Implementation of asynchronously send a record to a topic. Equivalent to <code>send(record, null)</code>. * See {@link #send(ProducerRecord, Callback)} for details. */ private Future<RecordMetadata> doSend(ProducerRecord<K, V> record, Callback callback) { TopicPartition tp = null; try { // first make sure the metadata for the topic is available long waitedOnMetadataMs = waitOnMetadata(record.topic(), this.maxBlockTimeMs); long remainingWaitMs = Math.max(0, this.maxBlockTimeMs - waitedOnMetadataMs); byte[] serializedKey; try { serializedKey = keySerializer.serialize(record.topic(), record.key()); } catch (ClassCastException cce) { throw new SerializationException("Can't convert key of class " + record.key().getClass().getName() + " to class " + producerConfig.getClass(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG).getName() + " specified in key.serializer"); } byte[] serializedValue; try { serializedValue = valueSerializer.serialize(record.topic(), record.value()); } catch (ClassCastException cce) { throw new SerializationException("Can't convert value of class " + record.value().getClass().getName() + " to class " + producerConfig.getClass(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG).getName() + " specified in value.serializer"); } int partition = partition(record, serializedKey, serializedValue, metadata.fetch()); int serializedSize = Records.LOG_OVERHEAD + Record.recordSize(serializedKey, serializedValue); ensureValidRecordSize(serializedSize); tp = new TopicPartition(record.topic(), partition); long timestamp = record.timestamp() == null ? time.milliseconds() : record.timestamp(); log.trace("Sending record {} with callback {} to topic {} partition {}", record, callback, record.topic(), partition); // producer callback will make sure to call both 'callback' and interceptor callback Callback interceptCallback = this.interceptors == null ? callback : new InterceptorCallback<>(callback, this.interceptors, tp); RecordAccumulator.RecordAppendResult result = accumulator.append(tp, timestamp, serializedKey, serializedValue, interceptCallback, remainingWaitMs); if (result.batchIsFull || result.newBatchCreated) { log.trace("Waking up the sender since topic {} partition {} is either full or getting a new batch", record.topic(), partition); this.sender.wakeup(); } return result.future; // handling exceptions and record the errors; // for API exceptions return them in the future, // for other exceptions throw directly } catch (ApiException e) { log.debug("Exception occurred during message send:", e); if (callback != null) callback.onCompletion(null, e); this.errors.record(); if (this.interceptors != null) this.interceptors.onSendError(record, tp, e); return new FutureFailure(e); } catch (InterruptedException e) { this.errors.record(); if (this.interceptors != null) this.interceptors.onSendError(record, tp, e); throw new InterruptException(e); } catch (BufferExhaustedException e) { this.errors.record(); this.metrics.sensor("buffer-exhausted-records").record(); if (this.interceptors != null) this.interceptors.onSendError(record, tp, e); throw e; } catch (KafkaException e) { this.errors.record(); if (this.interceptors != null) this.interceptors.onSendError(record, tp, e); throw e; } catch (Exception e) { // we notify interceptor about all exceptions, since onSend is called before anything else in this method if (this.interceptors != null) this.interceptors.onSendError(record, tp, e); throw e; } } 这个方法的脉络虽然比较长,但是脉络还是比较清晰,主要先执行了: 1)waitOnMetadata 应该是等待元数据拉取 2)keySerializer.serialize和valueSerializer.serialize,很明显就是将Record序列化成byte字节数组 3)通过partition进行路由分区,按照一定路由策略选择Topic下的某个分区 4)accumulator.append将消息放入缓冲器中 5)唤醒Sender线程的selector.select()的阻塞,开始处理内存缓冲器中的数据。 用图来表示如下所示: 这两节我们重点分析元数据拉取的这个场景的源码原理。 所以这里我们着重先看下步骤1 ,之后的4步我们之后会分析到的。 waitOnMetadata 如何等待元数据拉取的? 既然send的第一步是执行waitOnMetadata方法,首先看下它的代码: /** * Wait for cluster metadata including partitions for the given topic to be available. * @param topic The topic we want metadata for * @param maxWaitMs The maximum time in ms for waiting on the metadata * @return The amount of time we waited in ms */ private long waitOnMetadata(String topic, long maxWaitMs) throws InterruptedException { // add topic to metadata topic list if it is not there already. if (!this.metadata.containsTopic(topic)) this.metadata.add(topic); if (metadata.fetch().partitionsForTopic(topic) != null) return 0; long begin = time.milliseconds(); long remainingWaitMs = maxWaitMs; while (metadata.fetch().partitionsForTopic(topic) == null) { log.trace("Requesting metadata update for topic {}.", topic); int version = metadata.requestUpdate(); sender.wakeup(); metadata.awaitUpdate(version, remainingWaitMs); long elapsed = time.milliseconds() - begin; if (elapsed >= maxWaitMs) throw new TimeoutException("Failed to update metadata after " + maxWaitMs + " ms."); if (metadata.fetch().unauthorizedTopics().contains(topic)) throw new TopicAuthorizationException(topic); remainingWaitMs = maxWaitMs - elapsed; } return time.milliseconds() - begin; } /** * Get the current cluster info without blocking */ public synchronized Cluster fetch() { return this.cluster; } public synchronized int requestUpdate() { this.needUpdate = true; return this.version; } /** * Wait for metadata update until the current version is larger than the last version we know of */ public synchronized void awaitUpdate(final int lastVersion, final long maxWaitMs) throws InterruptedException { if (maxWaitMs < 0) { throw new IllegalArgumentException("Max time to wait for metadata updates should not be < 0 milli seconds"); } long begin = System.currentTimeMillis(); long remainingWaitMs = maxWaitMs; while (this.version <= lastVersion) { if (remainingWaitMs != 0) wait(remainingWaitMs); long elapsed = System.currentTimeMillis() - begin; if (elapsed >= maxWaitMs) throw new TimeoutException("Failed to update metadata after " + maxWaitMs + " ms."); remainingWaitMs = maxWaitMs - elapsed; } } 这个方法核心就是判断了是否有Cluster元数据信息,如果没有,进行了如下操作: 1)metadata.requestUpdate(); 更新了一个needUpdate标记,这个值会影响之前maybeUpdate的metadataTimeout的计算,可以让metadataTimeout为0 2)sender.wakeup();唤醒之前nioSelector.select()的阻塞,继续执行 3)metadata.awaitUpdate(version, remainingWaitMs); 主要进行了版本比较,如果不是最新版本,调用了Metadata.wait()方法(wait方法是每个Object都会有的方法,一般和notify或者notifyAll组合使用) 整个过程我直接用图给大家表示一下,如下所示: 整个图就是今天我们分析的关键结果了,**这里通过两种阻塞和唤醒机制,一个是NIO中Selector的select()和wakeUp(),一个是MetaData对象的wait()和notifyAll()机制。**所以这里要结合之前Sender线程的阻塞逻辑一起来理解。 是不是很有意思一种使用,这里没有用任何线程的join、sleep、wait、park、unpark、notify这些方法。 小结 最后我们简单小结下,这里一节我们主要分析了如下Producer的源码原理: 初始化KafkaProducer时并没有去拉取元数据,但是创建了Selector组件,启动了Sender线程,select阻塞等待请求响应。由于还没有发送任何请求,所以初始化时并没有去真正拉取元数据。 真正拉取元数据是在第一次send方法调用时,会唤醒唤醒Selector之前阻塞的select(),进入第二次while循环,从而发送拉取元数据请求,并且通过Obejct.wait的机制等待60s,等到从Broker拉取元数据成功后,才会继续执行真正的生产消息的请求,否则会报拉取元数据超时异常。 这一节我们只是看到了进行了wait如何等待元数据拉取。 而唤醒Selector的select之后应该会进入第二次while循环 那第二次while循环如何发送请求拉取元数据请求,并且在成功后notifyAll()进行唤醒操作的呢? 让我们下一节继续分析,大家敬请期待! 我们下一节见! 本文由博客一文多发平台 OpenWrite 发布!

资源下载

更多资源
Nacos

Nacos

Nacos /nɑ:kəʊs/ 是 Dynamic Naming and Configuration Service 的首字母简称,一个易于构建 AI Agent 应用的动态服务发现、配置管理和AI智能体管理平台。Nacos 致力于帮助您发现、配置和管理微服务及AI智能体应用。Nacos 提供了一组简单易用的特性集,帮助您快速实现动态服务发现、服务配置、服务元数据、流量管理。Nacos 帮助您更敏捷和容易地构建、交付和管理微服务平台。

Spring

Spring

Spring框架(Spring Framework)是由Rod Johnson于2002年提出的开源Java企业级应用框架,旨在通过使用JavaBean替代传统EJB实现方式降低企业级编程开发的复杂性。该框架基于简单性、可测试性和松耦合性设计理念,提供核心容器、应用上下文、数据访问集成等模块,支持整合Hibernate、Struts等第三方框架,其适用范围不仅限于服务器端开发,绝大多数Java应用均可从中受益。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册