首页 文章 精选 留言 我的

精选列表

搜索[MiMo模型],共10000篇文章
优秀的个人博客,低调大师

深入理解Java内存模型(五)——锁

锁的释放-获取建立的happens before 关系 锁是java并发编程中最重要的同步机制。锁除了让临界区互斥执行外,还可以让释放锁的线程向获取同一个锁的线程发送消息。 下面是锁释放-获取的示例代码: class MonitorExample { int a = 0; public synchronized void writer() { //1 a++; //2 } //3 int i = a; //5 public synchronized void reader() { //4 …… } //6 } 假设线程A执行writer()方法,随后线程B执行reader()方法。根据happens before规则,这个过程包含的happens before 关系可以分为两类: ●根据程序次序规则,1 happens before 2, 2 happens before 3; 4 happens before 5, 5 happens before 6。 ●根据监视器锁规则,3 happens before 4。 ●根据happens before 的传递性,2 happens before 5。 在上图中,每一个箭头链接的两个节点,代表了一个happens before 关系。黑色箭头表示程序顺序规则;橙色箭头表示监视器锁规则;蓝色箭头表示组合这些规则后提供的happens before保证。 上图表示在线程A释放了锁之后,随后线程B获取同一个锁。在上图中,2 happens before 5。因此,线程A在释放锁之前所有可见的共享变量,在线程B获取同一个锁之后,将立刻变得对B线程可见。 锁释放和获取的内存语义 当线程释放锁时,JMM会把该线程对应的本地内存中的共享变量刷新到主内存中。以上面的MonitorExample程序为例,A线程释放锁后,共享数据的状态 当线程获取锁时,JMM会把该线程对应的本地内存置为无效。从而使得被监视器保护的临界区代码必须要从主内存中去读取共享变量。下面是锁获取的状态 对比锁释放-获取的内存语义与volatile写-读的内存语义,可以看出:锁释放与volatile写有相同的内存语义;锁获取与volatile读有相同的内存语义。 下面对锁释放和锁获取的内存语义做个总结: ●线程A释放一个锁,实质上是线程A向接下来将要获取这个锁的某个线程发出了(线程A对共享变量所做修改的)消息。 ●线程B获取一个锁,实质上是线程B接收了之前某个线程发出的(在释放这个锁之前对共享变量所做修改的)消息。 ●线程A释放锁,随后线程B获取这个锁,这个过程实质上是线程A通过主内存向线程B发送消息。 锁内存语义的实现 本文将借助ReentrantLock的源代码,来分析锁内存语义的具体实现机制。 请看下面的示例代码: class ReentrantLockExample { int a = 0; ReentrantLock lock = new ReentrantLock(); public void writer() { lock.lock(); //获取锁 try { public void reader () { a++; } finally { lock.unlock(); //释放锁 } } lock.unlock(); //释放锁 lock.lock(); //获取锁 try { int i = a; …… } finally { } } } 在ReentrantLock中,调用lock()方法获取锁;调用unlock()方法释放锁。 ReentrantLock的实现依赖于java同步器框架AbstractQueuedSynchronizer(本文简称之为AQS)。AQS使用一个整型的volatile变量(命名为state)来维护同步状态,马上我们会看到,这个volatile变量是ReentrantLock内存语义实现的关键。 ReentrantLock分为公平锁和非公平锁,我们首先分析公平锁。 使用公平锁时,加锁方法lock()的方法调用轨迹如下: ●ReentrantLock : lock() ●FairSync : lock() ●AbstractQueuedSynchronizer : acquire(int arg) ●ReentrantLock : tryAcquire(int acquires) 在第4步真正开始加锁,下面是该方法的源代码: protected final boolean tryAcquire(int acquires) { final Thread current = Thread.currentThread(); if (c == 0) { int c = getState(); //获取锁的开始,首先读volatile变量state if (isFirst(current) && else if (current == getExclusiveOwnerThread()) { compareAndSetState(0, acquires)) { setExclusiveOwnerThread(current); return true; } } setState(nextc); int nextc = c + acquires; if (nextc < 0) throw new Error("Maximum lock count exceeded"); return true; } } return false; 从上面源代码中我们可以看出,加锁方法首先读volatile变量state。 在使用公平锁时,解锁方法unlock()的方法调用轨迹如下: ●ReentrantLock : unlock() ●AbstractQueuedSynchronizer : release(int arg) ●Sync : tryRelease(int releases) 在第3步真正开始释放锁,下面是该方法的源代码: protected final boolean tryRelease(int releases) { int c = getState() - releases; if (Thread.currentThread() != getExclusiveOwnerThread()) throw new IllegalMonitorStateException(); boolean free = false; if (c == 0) { return free; free = true; setExclusiveOwnerThread(null); } setState(c); //释放锁的最后,写volatile变量state } 从上面的源代码我们可以看出,在释放锁的最后写volatile变量state。 公平锁在释放锁的最后写volatile变量state;在获取锁时首先读这个volatile变量。根据volatile的happens-before规则,释放锁的线程在写volatile变量之前可见的共享变量,在获取锁的线程读取同一个volatile变量后将立即变的对获取锁的线程可见。 现在我们分析非公平锁的内存语义的实现。 非公平锁的释放和公平锁完全一样,所以这里仅仅分析非公平锁的获取。 使用公平锁时,加锁方法lock()的方法调用轨迹如下: ●ReentrantLock : lock() ●NonfairSync : lock() ●AbstractQueuedSynchronizer : compareAndSetState(int expect, int update) 在第3步真正开始加锁,下面是该方法的源代码: protected final boolean compareAndSetState(int expect, int update) { return unsafe.compareAndSwapInt(this, stateOffset, expect, update); } 该方法以原子操作的方式更新state变量,本文把java的compareAndSet()方法调用简称为CAS。JDK文档对该方法的说明如下:如果当前状态值等于预期值,则以原子方式将同步状态设置为给定的更新值。此操作具有 volatile 读和写的内存语义。 这里我们分别从编译器和处理器的角度来分析,CAS如何同时具有volatile读和volatile写的内存语义。 前文我们提到过,编译器不会对volatile读与volatile读后面的任意内存操作重排序;编译器不会对volatile写与volatile写前面的任意内存操作重排序。组合这两个条件,意味着为了同时实现volatile读和volatile写的内存语义,编译器不能对CAS与CAS前面和后面的任意内存操作重排序。 下面我们来分析在常见的intel x86处理器中,CAS是如何同时具有volatile读和volatile写的内存语义的。 下面是sun.misc.Unsafe类的compareAndSwapInt()方法的源代码: public final native boolean compareAndSwapInt(Object o, long offset, int expected, int x); 可以看到这是个本地方法调用。这个本地方法在openjdk中依次调用的c++代码为:unsafe.cpp,atomic.cpp和atomicwindowsx86.inline.hpp。这个本地方法的最终实现在openjdk的如下位置:openjdk-7-fcs-src-b147-27jun2011openjdkhotspotsrcoscpuwindowsx86vm atomicwindowsx86.inline.hpp(对应于windows操作系统,X86处理器)。下面是对应于intel x86处理器的源代码的片段: // Adding a lock prefix to an instruction on MP machine // VC++ doesn't like the lock prefix to be on a single line // so we can't insert a label after the lock prefix. // By emitting a lock prefix, we can define a label after it. #define LOCK_IF_MP(mp) __asm cmp mp, 0 __asm je L0 __asm _emit 0xF0 __asm L0: // alternative for InterlockedCompareExchange inline jint Atomic::cmpxchg (jint exchange_value, volatile jint* dest, jint compare_value) { int mp = os::is_MP(); __asm { mov edx, dest mov ecx, exchange_value } mov eax, compare_value LOCK_IF_MP(mp) cmpxchg dword ptr [edx], ecx } 如上面源代码所示,程序会根据当前处理器的类型来决定是否为cmpxchg指令添加lock前缀。如果程序是在多处理器上运行,就为cmpxchg指令加上lock前缀(lock cmpxchg)。反之,如果程序是在单处理器上运行,就省略lock前缀(单处理器自身会维护单处理器内的顺序一致性,不需要lock前缀提供的内存屏障效果)。 intel的手册对lock前缀的说明如下: ●确保对内存的读-改-写操作原子执行。在Pentium及Pentium之前的处理器中,带有lock前缀的指令在执行期间会锁住总线,使得其他处理器暂时无法通过总线访问内存。很显然,这会带来昂贵的开销。从Pentium 4,Intel Xeon及P6处理器开始,intel在原有总线锁的基础上做了一个很有意义的优化:如果要访问的内存区域(area of memory)在lock前缀指令执行期间已经在处理器内部的缓存中被锁定(即包含该内存区域的缓存行当前处于独占或以修改状态),并且该内存区域被完全包含在单个缓存行(cache line)中,那么处理器将直接执行该指令。由于在指令执行期间该缓存行会一直被锁定,其它处理器无法读/写该指令要访问的内存区域,因此能保证指令执行的原子性。这个操作过程叫做缓存锁定(cache locking),缓存锁定将大大降低lock前缀指令的执行开销,但是当多处理器之间的竞争程度很高或者指令访问的内存地址未对齐时,仍然会锁住总线。 ●禁止该指令与之前和之后的读和写指令重排序。 ●把写缓冲区中的所有数据刷新到内存中。 上面的第2点和第3点所具有的内存屏障效果,足以同时实现volatile读和volatile写的内存语义。 经过上面的这些分析,现在我们终于能明白为什么JDK文档说CAS同时具有volatile读和volatile写的内存语义了。 现在对公平锁和非公平锁的内存语义做个总结: ●公平锁和非公平锁释放时,最后都要写一个volatile变量state。 ●公平锁获取时,首先会去读这个volatile变量。 ●非公平锁获取时,首先会用CAS更新这个volatile变量,这个操作同时具有volatile读和volatile写的内存语义。 从本文对ReentrantLock的分析可以看出,锁释放-获取的内存语义的实现至少有下面两种方式: ●利用volatile变量的写-读所具有的内存语义。 ●利用CAS所附带的volatile读和volatile写的内存语义。 concurrent包的实现 由于java的CAS同时具有 volatile 读和volatile写的内存语义,因此Java线程之间的通信现在有了下面四种方式: ●A线程写volatile变量,随后B线程读这个volatile变量。 ●A线程写volatile变量,随后B线程用CAS更新这个volatile变量。 ●A线程用CAS更新一个volatile变量,随后B线程用CAS更新这个volatile变量。 ●A线程用CAS更新一个volatile变量,随后B线程读这个volatile变量。 Java的CAS会使用现代处理器上提供的高效机器级别原子指令,这些原子指令以原子方式对内存执行读-改-写操作,这是在多处理器中实现同步的关键(从本质上来说,能够支持原子性读-改-写指令的计算机器,是顺序计算图灵机的异步等价机器,因此任何现代的多处理器都会去支持某种能对内存执行原子性读-改-写操作的原子指令)。同时,volatile变量的读/写和CAS可以实现线程之间的通信。把这些特性整合在一起,就形成了整个concurrent包得以实现的基石。如果我们仔细分析concurrent包的源代码实现,会发现一个通用化的实现模式: ●首先,声明共享变量为volatile; ●然后,使用CAS的原子条件更新来实现线程之间的同步; ●同时,配合以volatile的读/写和CAS所具有的volatile读和写的内存语义来实现线程之间的通信。 原文发布时间为:2018-09-27 本文来自云栖社区合作伙伴“Java杂记”,了解相关信息可以关注“Java杂记”。

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

Python全栈 Web(Flask框架、模型、sqlalchemy)

Flask的文件传输: 如果大批量上传数据的时候(如:大文件) 就不能使用网页上传了 主要是由于http协议不支持 需要使用单独的上传工具(c/s版的) URL不存在参数上限的问题,HTTP协议规范也没有对URL长度进行限制。 这个限制是特定的浏览器及服务器对它的限制。 IE对URL长度的限制是2083字节(2K+35字节)对于其他浏览器, 如FireFox,Netscape等,则没有长度限制, 这个时候其限制取决于服务器的操作系统。即如果url太长, 服务器可能会因为安全方面的设置从而拒绝请求或者发生不完整的数据请求。 post 理论上讲是没有大小限制的,HTTP协议规范也没有进行大小限制, 但实际上post所能传递的数据量大小取决于服务器的设置和内存大小。 因为我们一般post的数据量很少超过MB的,所以我们很少能感觉的到post的数据量限制,

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

前后分离模型之封装 Api 调用

Ajax 和异步处理 调用 API 访问数据采用的 Ajax 方式,这是一个异步过程,异步过程最基本的处理方式是事件或回调,其实这两种处理方式实现原理差不多,都需要在调用异步过程的时候传入一个在异步过程结束的时候调用的接口。比如 jQuery Ajax 的 success 就是典型的回调参数。不过使用 jQuery 处理异步推荐使用 Promise 处理方式。 Promise 处理方式也是通过注册回调函数来完成的。jQuery 的 Promise 和 ES6 的标准 Promise 有点不一样,但在 then 上可以兼容,通常称为 thenable。jQuery 的 Promise 没有提供 .catch() 接口,但它自己定义的 .done()、.fail() 和 .always() 三个注册回调的方式也很有特色,用起来很方便,它是在事件的方式来注册的(即,可以注册多个同类型的处理函数,在该触发的时候都会触发)。 当然更直观的一点的处理方式是使用 ES2017 带来的 async/await 方式,可以用同步代码的形式来写异步代码,当然也有一些坑在里面。对于前端工程师来说,最大的坑就是有些浏览器不支持,需要进行转译,所以如果前端代码没有构建过程,一般还是就用 ES5 的语法兼容性好一些(jQuery 的 Promise 是支持 ES5 的,但是标准 Promise 要 ES6 以后才可以使用)。 关于 JavaScript 异步处理相关的内容可以参考 从小小题目逐步走进 JavaScript 异步调用 闲谈异步调用“扁平”化 从地狱到天堂,Node 回调向 async/await 转变 理解 JavaScript 的 async/await 从不用 try-catch 实现的 async/await 语法说错误处理 自己封装工具函数 在处理 Ajax 的过程中,虽然有现成的库(比如 jQuery.ajax,axios 等),它毕竟是为了通用目的设计的,在使用的时候仍然不免繁琐。而在项目中,对 Api 进行调用的过程几乎都大同小异。如果设计得当,就连错误处理的方式都会是一样的。因此,在项目内的 Ajax 调用其实可以进行进一步的封装,使之在项目内使用起来更方便。如果接口方式发生变化,修改起来也更容易。 比如,当前接口要求使用 POST 方法调用(暂不考虑 RESTful),参数必须包括 action,返回的数据以 JSON 方式提供,如果出错,只要不是服务器异常都会返回特定的 JSON 数据,包括一个不等于 0 的 code 和可选的 message 属性。 那么用 jQuery 写这么一个 Ajax 调用,大概是这样 const apiUrl = "http://api.some.com/"; jQuery .ajax(url, { type: "post", dataType: "json", data: { action: "login", username: "uname", password: "passwd" } }) .done(function(data) { if (data.code) { alert(data.message || "登录失败!"); } else { window.location.assign("home"); } }) .fail(function() { alert("服务器错误"); }); 初步封装 同一项目中,这样的 Ajax 调用,基本上只有 data 部分和 .done 回调中的 else 部分不同,所以进行一次封装会大大减少代码量,可以这样封装 function appAjax(action, params) { var deffered = $.Deferred(); jQuery .ajax(apiUrl, { type: "post", dataType: "json", data: $.extend({ action: action }, params) }) .done(function(data) { // 当 code 为 0 或省略时,表示没有错误, // 其它值表示错误代码 if (data.code) { if (data.message) { // 如果服务器返回了消息,那么向用户呈现消息 // resolve(null),表示不需要后续进行业务处理 alert(data.message); deffered.resolve(); } else { // 如果服务器没返回消息,那么把 data 丢给外面的业务处理 deferred.reject(data); } } else { // 正常返回数据的情况 deffered.resolve(data); } }) .fail(function() { // Ajax 调用失败,向用户呈现消息,同时不需要进行后续的业务处理 alert("服务器错误"); deffered.resolve(); }); return deferred.promise(); } 而业务层的调用就很简单了 appAjax("login", { username: "uname", password: "passwd" }).done(function(data) { if (data) { window.location.assign("home"); } }).fail(function() { alert("登录失败"); }); 更换 API 调用接口 上面的封装对调用接口和返回数据进行了统一处理,把大部分项目接口约定的内容都处理掉了,剩下在每次调用时需要处理的就是纯粹的业务。 现在项目组决定不用 jQuery 的 Ajax,而是采用 axios 来调用 API(axios 不见得就比 jQuery 好,这里只是举例),那么只需要修改一下 appAjax() 的实现即可。所有业务调用都不需要修改。 假设现在的目标环境仍然是 ES5,那么需要第三方 Promise 提供,这里拟用 Bluebird,兼容原生 Promise 接口(在 HTML 中引入,未直接出现在 JS 代码中)。 function appAjax(action, params) { var deffered = $.Deferred(); axios .post(apiUrl, { data: $.extend({ action: action }, params) }) .then(function(data) { ... }, function() { ... }); return deferred.promise(); } 这次的封装采用了 axios 来实现 Web Api 调用。但是为了保持原来的接口(jQuery Promise 对象有提供 .done()、.fail() 和 .always() 事件处理),appAjax 仍然不得不返回 jQuery Promise。这样,即使所有地方都不再需要使用 jQuery,这里仍然得用。 项目中应该用还是不用 jQuery?请阅读为什么要用原生 JavaScript 代替 jQuery? 去除 jQuery 就只在这里使用 jQuery 总让人感觉如芒在背,想把它去掉。有两个办法 修改所有业务中的调用,去掉 .done()、.fail() 和 .always(),改成 .then()。这一步工作量较大,但基本无痛,因为 jQuery Promise 本身支持 .then()。但是有一点需要特别注意,这一点稍后说明 自己写个适配器,兼容 jQuery Promise 的接口,工作量也不小,但关键是要充分测试,避免差错。 上面提到第 1 种方法中有一点需要特别注意,那就是 .then() 和 .done() 系列函数在处理方式上有所不同。.then() 是按 Promise 的特性设计的,它返回的是另一个 Promise 对象;而 .done() 系列函数是按事件机制实现的,返回的是原来的 Promise 对象。所以像下面这样的代码在修改时就要注意了 appAjax(url, params) .done(function(data) { console.log("第 1 处处理", data) }) .done(function(data) { console.log("第 2 处处理", data) }); // 第 1 处处理 {} // 第 2 处处理 {} 简单的把 .done() 改成 .then() 之后(注意不需要使用 Bluebird,因为 jQuery Promise 支持 .then()) appAjax(url, params) .then(function(data) { console.log("第 1 处处理", data); }) .then(function(data) { console.log("第 2 处处理", data); }); // 第 1 处处理 {} // 第 2 处处理 undefined 原因上面已经讲了,这里正确的处理方式是合并多个 done 的代码,或者在 .then() 处理函数中返回 data: appAjax(url, params) .then(function(data) { console.log("第 1 处处理", data); return data; }) .then(function(data) { console.log("第 2 处处理", data); }); 使用 Promise 接口改善设计 我们的 appAjax() 接口部分也可以设计成 Promise 实现,这是一个更通用的接口。既使用不用 ES2015+ 特性,也可以使用像 jQuery Promise 或 Bluebird 这样的三方库提供的 Promise。 function appAjax(action, params) { // axios 依赖于 Promise,ES5 中可以使用 Bluebird 提供的 Promise return axios .post(apiUrl, { data: $.extend({ action: action }, params) }) .then(function(data) { // 这里调整了判断顺序,会让代码看起来更简洁 if (!data.code) { return data; } if (!data.message) { throw data; } alert(data.message); }, function() { alert("服务器错误"); }); } 不过现在前端有构建工具,可以使用 ES2015+ 配置 Babel,也可以使用 TypeScript …… 总之,选择很多,写起来也很方便。那么在设计的时候就不用局限于 ES5 所支持的内容了。所以可以考虑用 Promise + async/await 来实现 async function appAjax(action, params) { // axios 依赖于 Promise,ES5 中可以使用 Bluebird 提供的 Promise const data = await axios .post(apiUrl, { data: $.extend({ action: action }, params) }) // 这里模拟一个包含错误消息的结果,以便后面统一处理错误 // 这样就不需要用 try ... catch 了 .catch(() => ({ code: -1, message: "服务器错误" })); if (!data.code) { return data; } if (!data.message) { throw data; } alert(data.message); } 上面代码中使用 .catch() 来避免 try ... catch ... 的技巧在从不用 try-catch 实现的 async/await 语法说错误处理中提到过。 当然业务层调用也可以使用 async/await(记得写在 async 函数中): const data = await appAjax("login", { username: "uname", password: "passwd" }).catch(() => { alert("登录失败"); }); if (data) { window.location.assign("home"); } 对于多次 .done() 的改造: const data = await appAjax(url, params); console.log("第 1 处处理", data); console.log("第 2 处处理", data); 小结 本文以封装 Ajax 调用为例,看似在讲述异步调用。但实际想告诉大家的东西是:如何将一个常用的功能封装起来,实现代码重用和更简洁的调用;以及在封装的过程中需要考虑的问题——向前和向后的兼容性,在做工具函数封装的时候,应该尽量避免和某个特定的工具特性绑定,向公共标准靠拢——不知大家是否有所体会。

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

go-hbase的Scan模型源码分析

git地址在这里:https://github.com/Lazyshot/go-hbase 这是一个使用go操作hbase的行为。 分析scan行为 如何使用scan看下面这个例子,伪代码如下: func scan(phone string, start time.Time, end time.Time) ([]Loc, error) { ... client := hbase.NewClient(zks, "/hbase") client.SetLogLevel("DEBUG") scan := client.Scan(table) scan.StartRow = []byte(phone + strconv.Itoa(int(end.Unix()))) scan.StopRow = []byte(phone + strconv.Itoa(int(start.Unix()))) var locs []Loc scan.Map(func(ret *hbase.ResultRow) { var loc Loc for _, v := range ret.Columns { switch v.ColumnName { case "lbs:phone": loc.Phone = v.Value.String() case "lbs:lat": loc.Lat = v.Value.String() ... } } locs = append(locs, loc) }) return locs, nil } 首先是NewClient, 返回的结构是hbase.Client, 这个结构代表的是与hbase服务端交互的客户端实体。 这里没有什么好看的,倒是有一点要注意,在NewClient的时候,里面的zkRootReginPath是写死的,就是说hbase在zk中的地址是固定的。当然这个也是默认的。 func NewClient(zkHosts []string, zkRoot string) *Client { cl := &Client{ zkHosts: zkHosts, zkRoot: zkRoot, zkRootRegionPath: "/meta-region-server", servers: make(map[string]*connection), cachedRegionLocations: make(map[string]map[string]*regionInfo), prefetched: make(map[string]bool), maxRetries: max_action_retries, } cl.initZk() return cl } 下面是client.Scan client.Scan 返回的是 func newScan(table []byte, client *Client) *Scan { return &Scan{ client: client, table: table, nextStartRow: nil, families: make([][]byte, 0), qualifiers: make([][][]byte, 0), numCached: 100, closed: false, timeRange: nil, } } scan结构: type Scan struct { client *Client id uint64 table []byte StartRow []byte StopRow []byte families [][]byte qualifiers [][][]byte nextStartRow []byte numCached int closed bool //for filters timeRange *TimeRange location *regionInfo server *connection } 设置了开始位置,结束位置,就可以进行Map操作了。 func (s *Scan) Map(f func(*ResultRow)) { for { results := s.next() if results == nil { break } for _, v := range results { f(v) if s.closed { return } } } } 这个map的参数是一个函数f,没有返回值。框架的行为就是一个大循环,不断调用s.next(),注意,这里s.next返回回来的result可能是由多条,然后把这个多条数据每条进行一次实际的函数调用。结束循环有两个方法,一个是next中再也取不到数据(数据已经取完了)。还有一个是s.closed呗设置为true。 s.next() func (s *Scan) next() []*ResultRow { startRow := s.nextStartRow if startRow == nil { startRow = s.StartRow } return s.getData(startRow) } 这里其实是把startRow不断往前推进,但是每次从startRow获取多少数据呢?需要看getData getData 最核心的流程如下: func (s *Scan) getData(nextStart []byte) []*ResultRow { ... server, location := s.getServerAndLocation(s.table, nextStart) req := &proto.ScanRequest{ Region: &proto.RegionSpecifier{ Type: proto.RegionSpecifier_REGION_NAME.Enum(), Value: []byte(location.name), }, NumberOfRows: pb.Uint32(uint32(s.numCached)), Scan: &proto.Scan{}, } ... cl := newCall(req) server.call(cl) ... select { case msg := <-cl.responseCh: return s.processResponse(msg) } } 这里看到有一个s.numCached, 我们猜测这个是用来指定一次call请求调用回多少条数据的。 看call函数 func newCall(request pb.Message) *call { var responseBuffer pb.Message var methodName string switch request.(type) { ... case *proto.ScanRequest: responseBuffer = &proto.ScanResponse{} methodName = "Scan" ... } return &call{ methodName: methodName, request: request, responseBuffer: responseBuffer, responseCh: make(chan pb.Message, 1), } } type call struct { id uint32 methodName string request pb.Message responseBuffer pb.Message responseCh chan pb.Message } 可以看出,这个call是一个有responseBuffer的实际调用者。 下面看server.Call 至于这里的server, 我们不看代码流程了,只需要知道最后他返回的是connection这么个结构 type connection struct { connstr string id int name string socket net.Conn in *inputStream calls map[int]*call callId *atomicCounter isMaster bool } 创建是使用函数newConnection调用 func newConnection(connstr string, isMaster bool) (*connection, error) { id := connectionIds.IncrAndGet() log.Debug("Connecting to server[id=%d] [%s]", id, connstr) socket, err := net.Dial("tcp", connstr) if err != nil { return nil, err } c := &connection{ connstr: connstr, id: id, name: fmt.Sprintf("connection(%s) id: %d", connstr, id), socket: socket, in: newInputStream(socket), calls: make(map[int]*call), callId: newAtomicCounter(), isMaster: isMaster, } err = c.init() if err != nil { return nil, err } log.Debug("Initiated connection [id=%d] [%s]", id, connstr) return c, nil } 好,那么实际上就是调用connection.call(request *call) func (c *connection) call(request *call) error { id := c.callId.IncrAndGet() rh := &proto.RequestHeader{ CallId: pb.Uint32(uint32(id)), MethodName: pb.String(request.methodName), RequestParam: pb.Bool(true), } request.setid(uint32(id)) bfrh := newOutputBuffer() err := bfrh.WritePBMessage(rh) ... bfr := newOutputBuffer() err = bfr.WritePBMessage(request.request) ... buf := newOutputBuffer() buf.writeDelimitedBuffers(bfrh, bfr) c.calls[id] = request n, err := c.socket.Write(buf.Bytes()) ... } 逻辑就是先把requestHeader压入,再压入request.request call只是完成了请求转换成byte传输到hbase服务端,在什么地方进行消息回收呢? 回到NewConnection的方法,里面有个connection.init() func (c *connection) init() error { err := c.writeHead() if err != nil { return err } err = c.writeConnectionHeader() if err != nil { return err } go c.processMessages() return nil } 这里go c.processMessage() func (c *connection) processMessages() { for { msgs := c.in.processData() if msgs == nil || len(msgs) == 0 || len(msgs[0]) == 0 { continue } var rh proto.ResponseHeader err := pb.Unmarshal(msgs[0], &rh) if err != nil { panic(err) } callId := rh.GetCallId() call, ok := c.calls[int(callId)] delete(c.calls, int(callId)) exception := rh.GetException() if exception != nil { call.complete(fmt.Errorf("Exception returned: %s\n%s", exception.GetExceptionClassName(), exception.GetStackTrace()), nil) } else if len(msgs) == 2 { call.complete(nil, msgs[1]) } } } 这里将它简化下: func (c *connection) processMessages() { for { msgs := c.in.processData() call.complete(nil, msgs[1]) } } c.in.processData 是在input_stream.go中 func (in *inputStream) processData() [][]byte { nBytesExpecting, err := in.readInt32() ... if nBytesExpecting > 0 { buf, err := in.readN(nBytesExpecting) if err != nil && err == io.EOF { panic("Unexpected closed socket") } payloads := in.processMessage(buf) if len(payloads) > 0 { return payloads } } return nil } 先读取出一个int值,这个int值判断后面还有多少个bytes,再将后面的bytes读取进入到buf中,进行input_stream的processMessage处理。 我们这里还看到并没有执行我们map中定义的匿名方法。只是把消息解析出来了而已。 call.complete func (c *call) complete(err error, response []byte) { ... err2 := pb.Unmarshal(response, c.responseBuffer) ... c.responseCh <- c.responseBuffer } 这个函数有用的也就这两句话把responseBuffer里面的内容通过管道传递给responseCh 这里就看到getData的时候,被堵塞的地方 select { case msg := <-cl.responseCh: return s.processResponse(msg) } 那么这里就有把获取到的responseCh的消息进行processResponse处理。 func (s *Scan) processResponse(response pb.Message) []*ResultRow { ... results := res.GetResults() n := len(results) ... s.closeScan(s.server, s.location, s.id) ... tbr := make([]*ResultRow, n) for i, v := range results { tbr[i] = newResultRow(v) } return tbr } 这个函数并没有什么特别的行为,只是进行ResultRow的组装。 好吧,这个包有个地方可以优化,这个go-hbase的scan的时候,numCached默认是100,这个对于hbase来说太小了,完全可以调整大点,到2000~10000之间,你会发现scan的性能提升杠杠的。 本文转自轩脉刃博客园博客,原文链接:http://www.cnblogs.com/yjf512/p/6076690.html,如需转载请自行联系原作者

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

OSI七层模型及对应协议

OSI七个层次的功能 物理层 为数据链路层提供物理连接,在其上串行传送比特流,即所传送数据的单位是比特。此外,该层中还具有确定连接设备的电气特性和物理特性等功能。 数据链路层 负责在网络节点间的线路上通过检测、流量控制和重发等手段,无差错地传送以帧为单位的数据。为做到这一点,在每一帧中必须同时带有同步、地址、差错控制及流量控制等控制信息。 网络层 为了将数据分组从源(源端系统)送到目的地(目标端系统),网络层的任务就是选择合适的路由和交换节点,使源的传输层传下来的分组信息能够正确无误地按照地址找到目的地,并交付给相应的传输层,即完成网络的寻址功能。 传输层 传输层是高低层之间衔接的接口层。数据传输的单位是报文,当报文较长时将它分割成若干分组,然后交给网络层进行传输。传输层是计算机网络协议分层中的最关键一层,该层以上各层将不再管理信息传输问题。 会话层 该层对传输的报文提供同步管理服务。在两个不同系统的互相通信的应用进程之间建立、组织和协调交互。例如,确定是双工还是半双工工作。 表示层 该层的主要任务是把所传送的数据的抽象语法变换为传送语法,即把不同计算机内部的不同表示形式转换成网络通信中的标准表示形式。此外,对传送的数据加密(或解密)、正文压缩(或还原)也是表示层的任务。 应用层 该层直接面向用户,是OSI中的最高层。它的主要任务是为用户提供应用的接口,即提供不同计算机间的文件传送、访问与管理,电子邮件的内容处理,不同计算机通过网络交互访问的虚拟终端功能等。 OSI中的层 功能 TCP/IP协议族 应用层 文件传输,电子邮件,文件服务,虚拟终端 TFTP,HTTP,SNMP,FTP,SMTP,DNS,Telnet 表示层 数据格式化,代码转换,数据加密 没有协议 会话层 解除或建立与别的接点的联系 没有协议 传输层 提供端对端的接口 TCP,UDP 网络层 为数据包选择路由 IP,ICMP,RIP,OSPF,BGP,IGMP 数据链路层 传输有地址的帧以及错误检测功能 SLIP,CSLIP,PPP,ARP,RARP,MTU 物理层 以二进制数据形式在物理媒体上传输数据 ISO2110,IEEE802。IEEE802.2 image.png 个人介绍: 高广超:多年一线互联网研发与架构设计经验,擅长设计与落地高可用、高性能、可扩展的互联网架构。 本文首发在 高广超的简书博客 转载请注明! 简书博客 头条号

资源下载

更多资源
Mario

Mario

马里奥是站在游戏界顶峰的超人气多面角色。马里奥靠吃蘑菇成长,特征是大鼻子、头戴帽子、身穿背带裤,还留着胡子。与他的双胞胎兄弟路易基一起,长年担任任天堂的招牌角色。

Spring

Spring

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

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

Sublime Text

Sublime Text

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

用户登录
用户注册