0%

Java程序性能优化

img

OverStack

默认xss

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
java -XX:+PrintFlagsFinal -version | grep -i 'stack'
intx CompilerThreadStackSize = 0 {pd product}
uintx GCDrainStackTargetSize = 64 {product}
bool JavaMonitorsInStackTrace = true {product}
uintx MarkStackSize = 4194304 {product}
uintx MarkStackSizeMax = 536870912 {product}
intx MaxJavaStackTraceDepth = 1024 {product}
bool OmitStackTraceInFastThrow = true {product}
intx OnStackReplacePercentage = 140 {pd product}
intx StackRedPages = 1 {pd product}
intx StackShadowPages = 20 {pd product}
bool StackTraceInThrowable = true {product}
intx StackYellowPages = 2 {pd product}
intx ThreadStackSize = 1024 {pd product}
bool UseOnStackReplacement = true {pd product}
intx VMThreadStackSize = 1024 {pd product}
java version "1.8.0_211"
Java(TM) SE Runtime Environment (build 1.8.0_211-b12)
Java HotSpot(TM) 64-Bit Server VM (build 25.211-b12, mixed mode)

https://docs.oracle.com/javase/8/docs/technotes/tools/unix/java.html

-Xss is translated in a VM flag named ThreadStackSize, 是-XX:ThreadStackSize的简写形式, 即线程栈的大小, 单位kb

-Xsssize

Sets the thread stack size (in bytes). Append the
letter k or K to indicate KB, m or M to indicate MB, g or G to
indicate GB. The default value depends on the platform:

  • Linux/ARM (32-bit): 320 KB
  • Linux/i386 (32-bit): 320 KB
  • Linux/x64 (64-bit): 1024 KB
  • OS X (64-bit): 1024 KB
  • Oracle Solaris/i386 (32-bit): 320 KB
  • Oracle Solaris/x64 (64-bit): 1024 KB

The following examples set the thread stack size to 1024 KB in different units:

-Xss1m
-Xss1024k
-Xss1048576

This option is equivalent to -XX:ThreadStackSize.

OverStack样例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
public class OverStackTest {
private int cnt = 0;

private void recursion() {
cnt++;
recursion();
}

private void testOverStack() {
try {
recursion();
} catch (Throwable e) {
System.out.println("deep of stack is " + cnt);
}
}

public static void main(String[] args) {
new OverStackTest().testOverStack();
}
}
1
2
3
4
5
6
7
> javac OverStackTest.java
> java -Xss256k OverStackTest
deep of stack is 1887
> java -Xss512k OverStackTest
deep of stack is 7428
> java -Xss670k OverStackTest
deep of stack is 7616
  • 标记复制: 新生代垃圾回收的通用算法
  • 标记压缩: 老年代垃圾回收

CMS收集器

三种操作:

  • 年轻代回收(暂停所有应用线程)
  • 启动并发线程回收老年带空间
  • 如有必要,full gc

http://jameswxx.iteye.com/blog/1041173

jstack

jstack命令的语法格式: jstack 。可以用jps查看java进程id。这里要注意的是:

  1. 不同的 JAVA虚机的线程 DUMP的创建方法和文件格式是不一样的, 不同的 JVM版本, dump信息也有差别。本文中, 只以 SUN的 hotspot JVM 5.0_06 为例。
  2. 在实际运行中, 往往一次 dump的信息, 还不足以确认问题。建议产生三次 dump信息, 如果每次 dump都指向同一个问题, 我们才确定问题的典型性。

线程分析

JVM 线程

在线程中, 有一些 JVM内部的后台线程, 来执行譬如垃圾回收, 或者低内存的检测等等任务, 这些线程往往在 JVM初始化的时候就存在, 如下所示:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
"Low Memory Detector" daemon prio=10 tid=0x081465f8 nid=0x7 runnable [0x00000000..0x00000000]  
​"CompilerThread0" daemon prio=10 tid=0x08143c58 nid=0x6 waiting on condition [0x00000000..0xfb5fd798]
​"Signal Dispatcher" daemon prio=10 tid=0x08142f08 nid=0x5 waiting on condition [0x00000000..0x00000000]
​"Finalizer" daemon prio=10 tid=0x08137ca0 nid=0x4 in Object.wait() [0xfbeed000..0xfbeeddb8] ​
​ at java.lang.Object.wait(Native Method)
​ ​ - waiting on <0xef600848> (a java.lang.ref.ReferenceQueue$Lock)
​ ​ at java.lang.ref.ReferenceQueue.remove(ReferenceQueue.java:116)
​ ​ - locked <0xef600848> (a java.lang.ref.ReferenceQueue$Lock)
​ ​ at java.lang.ref.ReferenceQueue.remove(ReferenceQueue.java:132)
​ ​ at java.lang.ref.Finalizer$FinalizerThread.run(Finalizer.java:159) ​
​"Reference Handler" daemon prio=10 tid=0x081370f0 nid=0x3 in Object.wait() [0xfbf4a000..0xfbf4aa38]
​ ​ at java.lang.Object.wait(Native Method)
​ ​ - waiting on <0xef600758> (a java.lang.ref.Reference$Lock)
​ ​ at java.lang.Object.wait(Object.java:474)
​ ​ at java.lang.ref.Reference$ReferenceHandler.run(Reference.java:116)
​ ​ - locked <0xef600758> (a java.lang.ref.Reference$Lock)
"VM Thread" prio=10 tid=0x08134878 nid=0x2 runnable
"VM Periodic Task Thread" prio=10 tid=0x08147768 nid=0x8 waiting on condition</span>

​ 我们更多的是要观察用户级别的线程, 如下所示:

1
2
3
4
5
6
"Thread-1" prio=10 tid=0x08223860 nid=0xa waiting on condition [0xef47a000..0xef47ac38]   ​
​ at java.lang.Thread.sleep(Native Method) ​
​ at testthread.MySleepingThread.method2(MySleepingThread.java:53)
​ ​ - locked <0xef63d600> (a testthread.MySleepingThread)
​ ​ at testthread.MySleepingThread.run(MySleepingThread.java:35) ​
​ at java.lang.Thread.run(Thread.java:595) </span>

我们能看到:
​ * 线程的状态: waiting on condition
​ * 线程的调用栈
​ * 线程的当前锁住的资源: <0xef63d600>

线程的状态分析

正如我们刚看到的那样, 线程的状态是一个重要的指标, 它会显示在线程 Stacktrace 的头一行结尾的地方。那么线程常见的有哪些状态呢?线程在什么样的情况下会进入这种状态呢?我们能从中发现什么线索?

Runnable

该状态表示线程具备所有运行条件, 在运行队列中准备操作系统的调度, 或者正在运行。

Wait on condition

该状态出现在线程等待某个条件的发生。具体是什么原因, 可以结合 stacktrace 来分析。

最常见的情况是线程在等待网络的读写, 比如当网络数据没有准备好读时, 线程处于这种等待状态, 而一旦有数据准备好读之后, 线程会重新激活, 读取并处理数据。
在 Java引入 NewIO之前, 对于每个网络连接, 都有一个对应的线程来处理网络的读写操作, 即使没有可读写的数据, 线程仍然阻塞在读写操作上, 这样有可能造成资源浪费, 而且给操作系统的线程调度也带来压力。在 NewIO里采用了新的机制, 编写的服务器程序的性能和可扩展性都得到提高。

如果发现有大量的线程都在处在 Wait on condition, 从线程 stack看, 正等待网络读写, 这可能是一个网络瓶颈的征兆。因为网络阻塞导致线程无法执行。一种情况是网络非常忙, 几 乎消耗了所有的带宽, 仍然有大量数据等待网络读 写;另一种情况也可能是网络空闲, 但由于路由等问题, 导致包无法正常的到达。所以要结合系统的一些性能观察工具来综合分析, 比如 netstat统计单位时间的发送包的数目, 如果很明显超过了所在网络带宽的限制 ; 观察 cpu的利用率, 如果系统态的 CPU时间, 相对于用户态的 CPU时间比例较高;如果程序运行在 Solaris 10平台上, 可以用 dtrace工具看系统调用的情况, 如果观察到 read/write的系统调用的次数或者运行时间遥遥领先;这些都指向由于网络带宽所限导致的网络瓶颈。另外一种出现 Wait on condition的常见情况是该线程在 sleep, 等待 sleep的时间到了时候, 将被唤醒。

Waiting for monitor entry 和 in Object.wait()


在多线程的 JAVA程序中, 实现线程之间的同步, 就要说说 Monitor。 Monitor是 Java中用以实现线程之间的互斥与协作的主要手段, 它可以看成是对象或者 Class的锁。每一个对象都有, 也仅有一个 monitor。每个 Monitor在某个时刻, 只能被一个线程拥有, 该线程就是 “Active Thread”, 而其它线程都是 “Waiting Thread”, 分别在两个队列 “ Entry Set”和 “Wait Set”里面等候。在 “Entry Set”中等待的线程状态是 “Waiting for monitor entry”, 而在 “Wait Set”中等待的线程状态是 “in Object.wait()”。

先看 “Entry Set”里面的线程。我们称被 synchronized保护起来的代码段为临界区。当一个线程申请进入临界区时, 它就进入了 “Entry Set”队列。对应的 code就像:

1
2
3
synchronized(obj) {
.........
}

这时有两种可能性:

  1. 该 monitor不被其它线程拥有, Entry Set里面也没有其它等待线程。本线程即成为相应类或者对象的 Monitor的 Owner, 执行临界区的代码
  2. 该 monitor被其它线程拥有, 本线程在 Entry Set队列中等待。*

在第一种情况下, 线程将处于 “Runnable”的状态, 而第二种情况下, 线程 DUMP会显示处于 “waiting for monitor entry”。如下所示:

1
2
3
4
5
"Thread-0" prio=10 tid=0x08222eb0 nid=0x9 waiting for monitor entry [0xf927b000..0xf927bdb8]  ​
at testthread.WaitThread.run(WaitThread.java:39) ​
- waiting to lock <0xef63bf08> (a java.lang.Object) ​
- locked <0xef63beb8> (a java.util.ArrayList) ​
at java.lang.Thread.run(Thread.java:595)


临界区的设置, 是为了保证其内部的代码执行的原子性和完整性。但是因为临界区在任何时间只允许线程串行通过, 这 和我们多线程的程序的初衷是相反的。 如果在多线程的程序中, 大量使用 synchronized, 或者不适当的使用了它, 会造成大量线程在临界区的入口等待, 造成系统的性能大幅下降。如果在线程 DUMP中发现了这个情况, 应该审查源码, 改进程序。

现在我们再来看现在线程为什么会进入 “Wait Set”。当线程获得了 Monitor, 进入了临界区之后, 如果发现线程继续运行的条件没有满足, 它则调用对象(一般就是被 synchronized 的对象)的 wait() 方法, 放弃了 Monitor, 进入 “Wait Set”队列。只有当别的线程在该对象上调用了 notify() 或者 notifyAll() , “ Wait Set”队列中线程才得到机会去竞争, 但是只有一个线程获得对象的 Monitor, 恢复到运行态。在 “Wait Set”中的线程, DUMP中表现为: in Object.wait(), 类似于:

1
2
3
4
5
6
7
"Thread-1" prio=10 tid=0x08223250 nid=0xa in Object.wait() [0xef47a000..0xef47aa38]  
​ at java.lang.Object.wait(Native Method) ​
​ - waiting on <0xef63beb8> (a java.util.ArrayList) ​
​ at java.lang.Object.wait(Object.java:474) ​
​ at testthread.MyWaitThread.run(MyWaitThread.java:40) ​
​ - locked <0xef63beb8> (a java.util.ArrayList) ​
​ at java.lang.Thread.run(Thread.java:595)

仔细观察上面的 DUMP信息, 你会发现它有以下两行:

  • locked <0xef63beb8> (a java.util.ArrayList)
  • waiting on <0xef63beb8> (a java.util.ArrayList)
    这里需要解释一下, 为什么先 lock了这个对象, 然后又 waiting on同一个对象呢?让我们看看这个线程对应的代码:

Java代码

1
2
3
4
5
synchronized(obj) {  
​ .........
​ obj.wait();
​ .........
}

线程的执行中, 先用 synchronized 获得了这个对象的 Monitor(对应于 locked <0xef63beb8> )。当执行到 obj.wait(), 线程即放弃了 Monitor的所有权, 进入 “wait set”队列(对应于 waiting on <0xef63beb8> )。
​ 往往在你的程序中, 会出现多个类似的线程, 他们都有相似的 DUMP信息。这也可能是正常的。比如, 在程序中, 有多个服务线程, 设计成从一个队列里面读取请求数据。这个队列就是 lock以及 waiting on的对象。当队列为空的时候, 这些线程都会在这个队列上等待, 直到队列有了数据, 这些线程被 Notify, 当然只有一个线程获得了 lock, 继续执行, 而其它线程继续等待。

JDK 5.0 的 lock

​ 上面我们提到如果 synchronized和 monitor机制运用不当, 可能会造成多线程程序的性能问题。在 JDK 5.0中, 引入了 Lock机制, 从而使开发者能更灵活的开发高性能的并发多线程程序, 可以替代以往 JDK中的 synchronized和 Monitor的 机制。但是, 要注意的是, 因为 Lock类只是一个普通类, JVM无从得知 Lock对象的占用情况, 所以在线程 DUMP中, 也不会包含关于 Lock的信息, 关于死锁等问题, 就不如用 synchronized的编程方式容易识别。

案例分析

  1. 死锁
    在多线程程序的编写中, 如果不适当的运用同步机制, 则有可能造成程序的死锁, 经常表现为程序的停顿, 或者不再响应用户的请求。比如在下面这个示例中, 是个较为典型的死锁情况:
1
2
3
4
5
6
7
8
9
10
"Thread-1" prio=5 tid=0x00acc490 nid=0xe50 waiting for monitor entry [0x02d3f000  
..0x02d3fd68]
at deadlockthreads.TestThread.run(TestThread.java:31)
- waiting to lock <0x22c19f18> (a java.lang.Object)
- locked <0x22c19f20> (a java.lang.Object)
"Thread-0" prio=5 tid=0x00accdb0 nid=0xdec waiting for monitor entry [0x02cff000
..0x02cff9e8]
at deadlockthreads.TestThread.run(TestThread.java:31)
- waiting to lock <0x22c19f20> (a java.lang.Object)
- locked <0x22c19f18> (a java.lang.Object)

在 JAVA 5中加强了对死锁的检测。线程 Dump中可以直接报告出 Java级别的死锁, 如下所示:

1
2
3
4
5
6
7
8
Found one Java-level deadlock:  
=============================
"Thread-1":
waiting to lock monitor 0x0003f334 (object 0x22c19f18, a java.lang.Object),
which is held by "Thread-0"
"Thread-0":
waiting to lock monitor 0x0003f314 (object 0x22c19f20, a java.lang.Object),
which is held by "Thread-1"
  1. 热锁
    ​ 热锁, 也往往是导致系统性能瓶颈的主要因素。其表现特征为, 由于多个线程对临界区, 或者锁的竞争, 可能出现:
  • 频繁的线程的上下文切换:从操作系统对线程的调度来看, 当 线程在等待资源而阻塞的时候, 操作系统会将之切换出来, 放到等待的队列, 当线程获得资源之后, 调度算法会将这个线程切换进去, 放到执行队列中。 * 大量的系统调用:因为线程的上下文切换, 以及热锁的竞争, 或 者临界区的频繁的进出, 都可能导致大量的系统调用。 * 大部分 CPU开销用在 “系统态 ”:线程上下文切换, 和系统调用, 都会导致 CPU在 “系统态 ”运行, 换而言之, 虽然系统很忙碌, 但是 CPU用在 “用户态 ”的比例较小, 应用程序得不到充分的 CPU资源。
  • 随着 CPU数目的增多, 系统的性能反而下降。因为 CPU数目多, 同 时运行的线程就越多, 可能就会造成更频繁的线程上下文切换和系统态的 CPU开销, 从而导致更糟糕的性能。 *
    上面的描述, 都是一个 scalability(可扩展性)很差的系统的表现。从整体的性能指标看, 由于线程热锁的存在, 程序的响应时间会变长, 吞吐量会降低。< /span>

    那么, 怎么去了解 “热锁 ”出现在什么地方呢?一个重要的方法还是结合操作系统的各种工具观察系统资源使用状况, 以及收集 Java线程的 DUMP信息, 看线程都阻塞在什么方法上, 了解原因, 才能找到对应的解决方法。

    我们曾经遇到过这样的例子, 程序运行时, 出现了以上指出的各种现象, 通过观察操作系统的资源使用统计信息, 以及线程 DUMP信息, 确定了程序中热锁的存在, 并发现大多数的线程状态都是 Waiting for monitor entry或者 Wait on monitor, 且是阻塞在压缩和解压缩的方法上。后来采用第三方的压缩包 javalib替代 JDK自带的压缩包后, 系统的性能提高了几倍。

1. 什么是Fork/Join框架

Fork/Join框架是Java7提供了的一个用于并行执行任务的框架, 是一个把大任务分割成若干个小任务,最终汇总每个小任务结果后得到大任务结果的框架。

我们再通过Fork和Join这两个单词来理解下Fork/Join框架,Fork就是把一个大任务切分为若干子任务并行的执行,Join就是合并这些子任务的执行结果,最后得到这个大任务的结果。比如计算1+2+。。+10000,可以分割成10个子任务,每个子任务分别对1000个数进行求和,最终汇总这10个子任务的结果。Fork/Join的运行流程图如下:

img

2. 工作窃取算法

工作窃取(work-stealing)算法是指某个线程从其他队列里窃取任务来执行。工作窃取的运行流程图如下:

img

那么为什么需要使用工作窃取算法呢?假如我们需要做一个比较大的任务,我们可以把这个任务分割为若干互不依赖的子任务,为了减少线程间的竞争,于是把这些子任务分别放到不同的队列里,并为每个队列创建一个单独的线程来执行队列里的任务,线程和队列一一对应,比如A线程负责处理A队列里的任务。但是有的线程会先把自己队列里的任务干完,而其他线程对应的队列里还有任务等待处理。干完活的线程与其等着,不如去帮其他线程干活,于是它就去其他线程的队列里窃取一个任务来执行。而在这时它们会访问同一个队列,所以为了减少窃取任务线程和被窃取任务线程之间的竞争,通常会使用双端队列,被窃取任务线程永远从双端队列的头部拿任务执行,而窃取任务的线程永远从双端队列的尾部拿任务执行。

工作窃取算法的优点是充分利用线程进行并行计算,并减少了线程间的竞争,其缺点是在某些情况下还是存在竞争,比如双端队列里只有一个任务时。并且消耗了更多的系统资源,比如创建多个线程和多个双端队列。

3. Fork/Join框架的介绍

我们已经很清楚Fork/Join框架的需求了,那么我们可以思考一下,如果让我们来设计一个Fork/Join框架,该如何设计?这个思考有助于你理解Fork/Join框架的设计。

第一步分割任务。首先我们需要有一个fork类来把大任务分割成子任务,有可能子任务还是很大,所以还需要不停的分割,直到分割出的子任务足够小。

第二步执行任务并合并结果。分割的子任务分别放在双端队列里,然后几个启动线程分别从双端队列里获取任务执行。子任务执行完的结果都统一放在一个队列里,启动一个线程从队列里拿数据,然后合并这些数据。

Fork/Join使用两个类来完成以上两件事情:

  • ForkJoinTask:我们要使用ForkJoin框架,必须首先创建一个ForkJoin任务。它提供在任务中执行fork()和join()操作的机制,通常情况下我们不需要直接继承ForkJoinTask类,而只需要继承它的子类,Fork/Join框架提供了以下两个子类:
    • RecursiveAction:用于没有返回结果的任务。
    • RecursiveTask :用于有返回结果的任务。
  • ForkJoinPool :ForkJoinTask需要通过ForkJoinPool来执行,任务分割出的子任务会添加到当前工作线程所维护的双端队列中,进入队列的头部。当一个工作线程的队列里暂时没有任务时,它会随机从其他工作线程的队列的尾部获取一个任务。

4. 使用Fork/Join框架

让我们通过一个简单的需求来使用下Fork/Join框架,需求是:计算1+2+3+4的结果。

使用Fork/Join框架首先要考虑到的是如何分割任务,如果我们希望每个子任务最多执行两个数的相加,那么我们设置分割的阈值是2,由于是4个数字相加,所以Fork/Join框架会把这个任务fork成两个子任务,子任务一负责计算1+2,子任务二负责计算3+4,然后再join两个子任务的结果。

因为是有结果的任务,所以必须继承RecursiveTask,实现代码如下:

img

img

通过这个例子让我们再来进一步了解ForkJoinTask,ForkJoinTask与一般的任务的主要区别在于它需要实现compute方法,在这个方法里,首先需要判断任务是否足够小,如果足够小就直接执行任务。如果不足够小,就必须分割成两个子任务,每个子任务在调用fork方法时,又会进入compute方法,看看当前子任务是否需要继续分割成孙任务,如果不需要继续分割,则执行当前子任务并返回结果。使用join方法会等待子任务执行完并得到其结果。

5. Fork/Join框架的异常处理

ForkJoinTask在执行的时候可能会抛出异常,但是我们没办法在主线程里直接捕获异常,所以ForkJoinTask提供了isCompletedAbnormally()方法来检查任务是否已经抛出异常或已经被取消了,并且可以通过ForkJoinTask的getException方法获取异常。使用如下代码:

1
2
3
4
5
if(task.isCompletedAbnormally())
{
System.out.println(task.getException());
}

getException方法返回Throwable对象,如果任务被取消了则返回CancellationException。如果任务没有完成或者没有抛出异常则返回null。

6. Fork/Join框架的实现原理

ForkJoinPool由ForkJoinTask数组和ForkJoinWorkerThread数组组成,ForkJoinTask数组负责存放程序提交给ForkJoinPool的任务,而ForkJoinWorkerThread数组负责执行这些任务。

ForkJoinTask的fork方法实现原理。当我们调用ForkJoinTask的fork方法时,程序会调用ForkJoinWorkerThread的pushTask方法异步的执行这个任务,然后立即返回结果。代码如下:

1
public final ForkJoinTask fork() {         ((ForkJoinWorkerThread) Thread.currentThread())             .pushTask(this);         return this; } 

pushTask方法把当前任务存放在ForkJoinTask 数组queue里。然后再调用ForkJoinPool的signalWork()方法唤醒或创建一个工作线程来执行任务。代码如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
final void pushTask(ForkJoinTask t) {
ForkJoinTask[] q; int s, m;
if ((q = queue) != null) { // ignore if queue removed
long u = (((s = queueTop) & (m = q.length - 1)) << ASHIFT) + ABASE;
UNSAFE.putOrderedObject(q, u, t);
queueTop = s + 1; // or use putOrderedInt
if ((s -= queueBase) <= 2)
pool.signalWork();
else if (s == m)
growQueue();
}
}

ForkJoinTask的join方法实现原理。Join方法的主要作用是阻塞当前线程并等待获取结果。让我们一起看看ForkJoinTask的join方法的实现,代码如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public final V join() {
if (doJoin() != NORMAL)
return reportResult();
else
return getRawResult();
}
private V reportResult() {
int s; Throwable ex;
if ((s = status) == CANCELLED)
throw new CancellationException();
if (s == EXCEPTIONAL && (ex = getThrowableException()) != null)
UNSAFE.throwException(ex);
return getRawResult();
}

首先,它调用了doJoin()方法,通过doJoin()方法得到当前任务的状态来判断返回什么结果,任务状态有四种:已完成(NORMAL),被取消(CANCELLED),信号(SIGNAL)和出现异常(EXCEPTIONAL)。

  • 如果任务状态是已完成,则直接返回任务结果。
  • 如果任务状态是被取消,则直接抛出CancellationException。
  • 如果任务状态是抛出异常,则直接抛出对应的异常。

让我们再来分析下doJoin()方法的实现代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
private int doJoin() {
Thread t; ForkJoinWorkerThread w; int s; boolean completed;
if ((t = Thread.currentThread()) instanceof ForkJoinWorkerThread) {
if ((s = status) < 0)
return s;
if ((w = (ForkJoinWorkerThread)t).unpushTask(this)) {
try {
completed = exec();
} catch (Throwable rex) {
return setExceptionalCompletion(rex);
}
if (completed)
return setCompletion(NORMAL);
}
return w.joinTask(this);
}
else
return externalAwaitDone();
}

在doJoin()方法里,首先通过查看任务的状态,看任务是否已经执行完了,如果执行完了,则直接返回任务状态,如果没有执行完,则从任务数组里取出任务并执行。如果任务顺利执行完成了,则设置任务状态为NORMAL,如果出现异常,则纪录异常,并将任务状态设置为EXCEPTIONAL。

7. 参考资料

8. 作者介绍

方腾飞,花名清英,并发编程网站站长。目前在阿里巴巴微贷事业部工作。并发编程网:http://ifeve.com,个人微博:http://weibo.com/kirals,欢迎通过我的微博进行技术交流。

感谢张龙对本文的审校。

网易镜像 https://c.163yun.com/hub#/m/home/

修改数据源

1
2
chmod a+w /etc/sysconfig/docker
## 增加 ADD_REGISTRY='--add-registry hub.c.163.com'

基本命令

image

1
docker image ls

docker

1
2
3
4
5
docker container ls -a
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
f7c2d84561dd hub.c.163.com/public/centos "/usr/sbin/sshd -D" 7 minutes ago Up 7 minutes 0.0.0.0:32768->22/tcp centos

docker exec -it centos /bin/bash

文件拷贝

  1. 本地copy文件到container
1
docker cp ~/putMerge.jar hadoop0:~/
  1. container文件copy到本地
1
docker cp hadoop0:~/putMerge.jar ~/
  1. images 重命名
1
2
3
4
docker tag IMAGEID(镜像id) REPOSITORY:TAG(仓库:标签)

#例子
docker tag ca1b6b825289 reffs/xxxxxxx:v1.0
1
2
3
4
5
6
7
8
9
10
yum list installed | grep docker
卸载:

yum -y remove docker.x86_64
重装:

yum install docker-io
启动:

service docker start

Docker 核心基础技术

  • Namespace
  • cgroup

linux Namespace

1
setns(fd,..)
1
2
3
4
5
6
7
8
9
10
ls -l /proc/1/ns

ll /proc/1/ns/
total 0
lrwxrwxrwx 1 root root 0 Jun 5 15:20 cgroup -> cgroup:[4026531835]
lrwxrwxrwx 1 root root 0 Jun 5 15:20 ipc -> ipc:[4026531839]
lrwxrwxrwx 1 root root 0 Apr 15 19:28 mnt -> mnt:[4026531840]
lrwxrwxrwx 1 root root 0 Apr 15 19:28 net -> net:[4026531957]
lrwxrwxrwx 1 root root 0 Jun 5 15:20 pid -> pid:[4026531836]
lrwxrwxrwx 1 root root 0 Apr 15 16:20 uts -> uts:[4026531838]

cgroup 子系统

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
lssubsys -m

lssubsys -m
cpuset /sys/fs/cgroup/cpuset
cpu,cpuacct /sys/fs/cgroup/cpu,cpuacct
memory /sys/fs/cgroup/memory
devices /sys/fs/cgroup/devices
freezer /sys/fs/cgroup/freezer
net_cls /sys/fs/cgroup/net_cls
blkio /sys/fs/cgroup/blkio
perf_event /sys/fs/cgroup/perf_event
hugetlb /sys/fs/cgroup/hugetlb


ls /sys/fs/cgroup/cpu/mytest

cgcreate -g cpu:mytest

M1芯片兼容

1
docker build --platform linux/amd64 -t image-name .

国内加速

macOS在docker-desktop的setting中添加:

1
2
3
4
5
6
"registry-mirrors": [
"https://mirror.ccs.tencentyun.com",
"https://dockerhub.woa.com",
"https://hub-mirror.c.163.com",
"https://mirror.baidubce.com"
]

yarn-2.8

yarn 2.8 访问 hbase 时,会带用户,hbase需要修改配置

Yarn 2.8

Spark 写入 SSD HBase,配置后可以正常写入

yarn.resourcemanager.principal gaiaadmin

mapReducer 处理数据的流程

Map Reducer shuffle

partition sorting combiner

shuffle过程

MapReducer过程

参考文献: http://www.cnblogs.com/ljy2013/articles/4435657.html

MapReducer 实践

  1. wordCount解析
  2. Mapper class, CombinerClass,ReducerClass
  3. 分析 TokenCountMapper 与 LongSumReducer
  4. 常用的系列化类: 实现 WriteComparable 的类

文件的读写

InputFormat

  • TextInputFormat 是 InputFormat 的默认实现, TextInputFormat
    返回的键是每行的字节偏移量,返回的值是该行的数据。
    key: LongWritable
    value: Text

  • KeyValueTextInputFormat 使用分隔符分割每行,分隔符之前的是键,之后的是值。默认的分隔符是制表符(\T),分离器的属性通过
    key.value.separator.in.input.line中指定

key: Text
value: Text

  • SequenceFileInputFormat<K,V> 用户自定义的序列化格式, 序列化文件为Hadoop专用的压缩二进制文件格式

key: K 用户定义
value: V 用户定义

  • NLineInputFormat key为分片的偏移量,value为包含N行数据的片段,N通过属性 mapred.line.inout.format.linespermap中指定,默认为1
    key: LongWritable
    value: Text

通过 JobConf.setInputFormat(KeyValueTextInputFormat.class)指定

自定义 InputFormat

1
2
3
4
5
6
7
8
9
10
11
12
public interface InputFormat<K,V> {
/**
* 将输入的文件分割成 numSplits 个片段
*/
InputSplit[] getSplits(JobConf job, int numSplits) throws IOException;
/**
*
*/
RecordReader<K,V> getRecordReader( InputSplit split,
JobConf job,
Reporter reporter) throws IOException;
}

上述所有的 InputFormat 都是 FileInputFormat 类的一个子类,FileInputFormat默认实现了 InputFormat 接口,且实现了 getSplits 方法,把输入的数据粗略的划分为一组分片, 每个分片的大小必须大于mapred.min.split.size个字节,且小于文件系统的块, 在实际情况下,一个分片的大小总是一个块的大小,在HDFS中默认为64MB; 而 getRecordReader为抽象方法, 在自定义实现InputFormat时, 可以通过继承 FileInputFormat , 自定义 getRecordReader 方法.

  • isSplitable 方法

  • RecordReader

  • 编写 InputFormat 实例

OutputFormat


【参考文献】

  1. ref-title

https://cwiki.apache.org/confluence/display/Hive/LanguageManual+ORC

1 文件格式

ORC的全称是(Optimized Row Columnar),它并不是一个单纯的列式存储格式,仍然是首先根据行组分割整个表,在每一个行组内进行按列存储。ORC文件是自描述的,它的元数据使用Protocol Buffers序列化,并且文件中的数据尽可能的压缩以降低存储空间的消耗。

1.1 列存储的优势

  • 查询的时候不需要扫描全部的数据,而只需要读取每次查询涉及的列,这样可以将I/O消耗降低N倍,另外可以保存每一列的统计信息(min、max、sum等),实现部分的谓词下推。
  • 由于每一列的成员都是同构的,可以针对不同的数据类型使用更高效的数据压缩算法,进一步减小I/O。
  • 由于每一列的成员的同构性,可以使用更加适合CPU pipeline的编码方式,减小CPU的缓存失效。

1.2 文件结构

ORC文件也是以二进制方式存储的,ORC文件也是自解析的,它包含许多的元数据,这些元数据都是同构ProtoBuffer进行序列化的。

  • ORC文件:保存在文件系统上的普通二进制文件,一个ORC文件中可以包含多个stripe,每一个stripe包含多条记录,这些记录按照列进行独立存储,对应到Parquet中的row group的概念。
  • 文件级元数据:包括文件的描述信息PostScript、文件meta信息(包括整个文件的统计信息)、所有stripe的信息和文件schema信息。
  • stripe:一组行形成一个stripe,每次读取文件是以行组为单位的,一般为HDFS的块大小,保存了每一列的索引和数据。
  • stripe元数据:保存stripe的位置、每一个列的在该stripe的统计信息以及所有的stream类型和位置。
  • row group:索引的最小单位,一个stripe中包含多个row group,默认为10000个值组成。
  • stream:一个stream表示文件中一段有效的数据,包括索引和数据两类。索引stream保存每一个row group的位置和统计信息,数据stream包括多种类型的数据,具体需要哪几种是由该列类型和编码方式决定。

1.3 基于统计信息的过滤

三个层级的统计信息,分别为文件级别、stripe级别和row group级别的,都可以用来根据Search ARGuments(谓词下推条件)判断是否可以跳过某些数据,在统计信息中都包含成员数和是否有null值,并且对于不同类型的数据设置一些特定的统计信息。

(1)file level 在ORC文件的末尾会记录文件级别的统计信息,会记录整个文件中columns的统计信息。这些信息主要用于查询的优化,也可以为一些简单的聚合查询比如max, min, sum输出结果。 
(2)stripe level ORC文件会保存每个字段stripe级别的统计信息,ORC reader使用这些统计信息来确定对于一个查询语句来说,需要读入哪些stripe中的记录。比如说某个stripe的字段max(a)=10,min(a)=3,那么当where条件为a >10或者a <3时,那么这个stripe中的所有记录在查询语句执行时不会被读入。 
(3)row level 为了进一步的避免读入不必要的数据,在逻辑上将一个column的index以一个给定的值(默认为10000,可由参数配置)分割为多个index组。以10000条记录为一个组,对数据进行统计。Hive查询引擎会将where条件中的约束传递给ORC reader,这些reader根据组级别的统计信息,过滤掉不必要的数据。如果该值设置的太小,就会保存更多的统计信息,用户需要根据自己数据的特点权衡一个合理的值。

1.4 Stripe → Row Group → Column → Stream 的关系

1
2
3
4
5
6
7
8
9
10
11
12
Stripe
├── Row Group 0
│ ├── Column A
│ │ ├── Stream PRESENT
│ │ ├── Stream DATA
│ │ └── Stream LENGTH / DICTIONARY / ...
│ ├── Column B
│ │ ├── Stream PRESENT
│ │ └── Stream DATA
│ └── ...
├── Row Group 1
└── ...
  • Row Group 是索引与统计的最小过滤单位
  • 真正的数据存储是按 column + stream 来的
  • Row Group 并不单独存数据块,而是通过 ROW_INDEX 指向各个 stream 在文件中的偏移位置
  • 每个Stream只属于一个字段,但一个字段可以包含多个Stream

2 ORC file API

源码主要有core和mapreduce两个模块,其中MapReduce 模块就是Hadoop的InputFormat和OutFormat,这个模块下有两个包,mapred和MapReduce分别符合MapReduce的V1和V2的Input/Output Format.

四个类:OrcInputFormat、OrcOutputFormat、OrcMapreduceRecordReader和OrcMapreduceRecordWriter

2.1 OrcInputFormat与OrcOutputFormat

OrcInputFormat继承了FileInputFormat,getSplits用的就是FileInputFormat的实现,只是重写了createRecordReader方法。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
@Override
public RecordReader<NullWritable, V>
createRecordReader(InputSplit inputSplit,
TaskAttemptContext taskAttemptContext
) throws IOException, InterruptedException {
FileSplit split = (FileSplit) inputSplit;
Configuration conf = taskAttemptContext.getConfiguration();
Reader file = OrcFile.createReader(split.getPath(),
OrcFile.readerOptions(conf)
.maxLength(OrcConf.MAX_FILE_LENGTH.getLong(conf)));
return new OrcMapreduceRecordReader<>(file,
org.apache.orc.mapred.OrcInputFormat.buildOptions(conf,
file, split.getStart(), split.getLength()));
}

这里面用到了org.apache.orc.Reader和org.apache.orc.OrcFile. 模板里面的V通常会是OrcStruct,Orc中一个表的schema可以表示为一个OrcStruct。
同样滴,OrcOutputFormat当中也是重写了craeteRecordWriter方法:

1
2
3
4
5
6
7
8
9
10
@Override
public RecordWriter<NullWritable, V>
getRecordWriter(TaskAttemptContext taskAttemptContext
) throws IOException {
Configuration conf = taskAttemptContext.getConfiguration();
Path filename = getDefaultWorkFile(taskAttemptContext, EXTENSION);
Writer writer = OrcFile.createWriter(filename,
org.apache.orc.mapred.OrcOutputFormat.buildOptions(conf));
return new OrcMapreduceRecordWriter<V>(writer);
}

这里面用到了org.apache.orc.Writer和org.apache.orc.OrcFile.
以上用到的这三个类都是core 模块里的,之后再去读。

2.2 OrcMapreduceRecordReader

这个类里面其实复用了org.apache.orc.mapred包下面的一些Writable类,这些Writable类用来在MapReduce中支持Orc中的一些特殊数据类型,比如Map、List、Struct等等,可以将这些数据类型从DataInput中反序列化出来,也可以将数据序列化到DataOutput中去。

从构造方法可以看到:

1
2
3
4
5
6
7
8
9
10
11
12
public OrcMapreduceRecordReader(Reader fileReader,
Reader.Options options) throws IOException {
this.batchReader = fileReader.rows(options);
if (options.getSchema() == null) {
schema = fileReader.getSchema();
} else {
schema = options.getSchema();
}
this.batch = schema.createRowBatch();
rowInBatch = 0;
this.row = (V) OrcStruct.createValue(schema);
}

其中batchReader是org.apache.orc.RecordReader类型的,这个RecordReader在core中。batch是真正有数据的地方。batch的类型是org.apache.hadoop.hive.ql.exec.vector.VectorizedRowBatch,从java doc可以看出来它是干嘛用的:

1
2
3
4
5
6
7
/**
* A VectorizedRowBatch is a set of rows, organized with each column
* as a vector. It is the unit of query execution, organized to minimize
* the cost per row and achieve high cycles-per-instruction.
* The major fields are public by design to allow fast and convenient
* access by the vectorized query execution code.
*/

而schema是org.apache.orc.TypeDescription类型的,而TypeDescription其实是一个数据类型的描述,它有一个category,可以是INT、FLOAT等等,也可以是STRUCT。这里schema的category通常会是STRUCT。STRUCT和OrcStruct是对应的,是一个复合类型,其中可以包含其他类型的属性,用来表示一个表的schema。

row就是batchReader从batch中读出来的一行记录。

OrcMapreduceRecordReader中另外一个重要的方法是nextKeyValue:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
@Override
public boolean nextKeyValue() throws IOException, InterruptedException {
if (!ensureBatch()) {
return false;
}
if (schema.getCategory() == TypeDescription.Category.STRUCT) {
OrcStruct result = (OrcStruct) row;
List<TypeDescription> children = schema.getChildren();
int numberOfChildren = children.size();
for(int i=0; i < numberOfChildren; ++i) {
result.setFieldValue(i, OrcMapredRecordReader.nextValue(batch.cols[i], rowInBatch,
children.get(i), result.getFieldValue(i)));
}
} else {
OrcMapredRecordReader.nextValue(batch.cols[0], rowInBatch, schema, row);
}
rowInBatch += 1;
return true;
}

可以看出来,这个方法中,判断当schema是STRUCT的category时,将它看做一个表的schema,取出里面包含的children,即各个字段的TypeDescription,然后读取各个字段的值,存入row中。这其中复用了OrcMapredRecordReader.nextValue()方法.

2.3 OrcMapreduceRecordWriter

OrcMapreduceRecordWriter和OrcMapreduceRecordReader类似,只不过其中的batchReader换成了org.apache.orc
.Writer类型的writer。构造方法如下:

1
2
3
4
5
6
public OrcMapreduceRecordWriter(Writer writer) {
this.writer = writer;
schema = writer.getSchema();
this.batch = schema.createRowBatch();
isTopStruct = schema.getCategory() == TypeDescription.Category.STRUCT;
}

这个类的主要方法就是write方法:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
@Override
public void write(NullWritable nullWritable, V v) throws IOException {
// if the batch is full, write it out.
if (batch.size == batch.getMaxSize()) {
writer.addRowBatch(batch);
batch.reset();
}

// add the new row
int row = batch.size++;
// skip over the OrcKey or OrcValue
if (v instanceof OrcKey) {
v = (V)((OrcKey) v).key;
} else if (v instanceof OrcValue) {
v = (V)((OrcValue) v).value;
}
if (isTopStruct) {
for(int f=0; f < schema.getChildren().size(); ++f) {
OrcMapredRecordWriter.setColumn(schema.getChildren().get(f),
batch.cols[f], row, ((OrcStruct) v).getFieldValue(f));
}
} else {
OrcMapredRecordWriter.setColumn(schema, batch.cols[0], row, v);
}
}

可以看出来,write先把V类型(通常是OrcStruct)的记录写入batch,如果batch写满了,就批量写入writer。

2.4 代办

  • Flink消费orc文件,能否按照file+strip作为offset?

2.5 【参考文献】

  1. ORC源码阅读(1) - mapreduce 模块

Hadoop是采用Java实现的,所有的命令和操作全部采用Java完成。
HDFS: Hadoop Distributed File System 是Hadoop提供的分布式文件系统。

本文介绍HDFS的Java API

1
2
3
4
5
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-core</artifactId>
<version>1.2.1</version>
</dependency>

首先,看一个例子:

putMerge: 本地文件合并并保存到HDFS

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
/**
* copy files in <code>srcDir</code>
* and merge to <code>objectFileName</code> in HDFS
*
* @param srcDir 本地文件夹
* @param objectFileName HDFS文件名
*/
private static void putMerge(String srcDir, String objectFileName) {
System.out.println("srcDir:" + srcDir+"\t objDir:" + objectFileName);
// 读取HDFS的默认配置, 配置文件主要包括:
// core-default.xml, core-site.xml, mapred-default.xml, mapred-site.xml, yarn-default.xml, yarn-site.xml, hdfs-default.xml, hdfs-site.xml
Configuration conf = new Configuration();
FileSystem hdfs = null;
LocalFileSystem local = null;

try {
hdfs = FileSystem.get(conf);// 获取HDFS配置的FileSystem
local = FileSystem.getLocal(conf);// 获取本地配置FileSystem
Path inputDir = new Path(srcDir); // 设定输入目录
System.out.println("inputDir: " + inputDir);
System.out.println("local.homeDirectory: " + local.getHomeDirectory());
System.out.println("local.workingDirectory: " + local.getWorkingDirectory());
System.out.println("local.uri: " + local.getUri());
System.out.println("local.conf: " + local.getConf());

FileStatus[] inputFiles = local.listStatus(inputDir);
System.out.println("inputFiles: " + inputFiles);
for (FileStatus inputFile : inputFiles) {
System.out.println("inputFile: " + inputFile);
}
if (inputFiles.length == 0) {
System.err.println("input file path is empty!");
System.exit(1);
}
Path hdfsFile = new Path(objectFileName); // 设定输出文件名称
FSDataOutputStream out = hdfs.create(hdfsFile, new Progressable() {
@Override
public void progress() {
System.out.print("*");
}
});
for (FileStatus inputFile : inputFiles) {
System.out.println("inputFile.path.name: " + inputFile.getPath().getName());
FSDataInputStream in = local.open(inputFile.getPath());
byte buffer[] = new byte[256];
int bytesRead = 0;
while ((bytesRead = in.read(buffer)) > 0) {
out.write(buffer, 0, bytesRead);
}
in.close();
System.out.println();
}
out.close();
} catch (IOException e) {
e.printStackTrace();
}
System.out.println();
}

执行:

  1. 打包
1
mvn package
  1. 复制到container
1
docker cp ~/IdeaProjects/MyTest/out/artifacts/HadoopTest_jar/HadoopTest.jar hadoop0:/root/putMerge.jar
  1. 执行
1
hadoop jar putMerge.jar en002 en002.txt

执行结果

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
version 1
args:en002
args:en002.txt
srcDir:en002
objDir:en002.txt
inputDir: en002
local.homeDirectory: file:/root
local.workingDirectory: file:/root
local.uri: file:///
local.conf: Configuration: core-default.xml, core-site.xml, mapred-default.xml, mapred-site.xml, yarn-default.xml, yarn-site.xml, hdfs-default.xml, hdfs-site.xml
inputFiles: [Lorg.apache.hadoop.fs.FileStatus;@4073c6c9
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/shengjing.txt; isDirectory=false; length=4467663; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/at.txt; isDirectory=false; length=829203; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/abc.txt; isDirectory=false; length=0; replication=1; blocksize=33554432; modification_time=1495543277000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/a.txt; isDirectory=false; length=18516; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/av.txt; isDirectory=false; length=189407; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/David.txt; isDirectory=false; length=1519616; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/Oliver.txt; isDirectory=false; length=981553; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/Jane.txt; isDirectory=false; length=1114997; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile: DeprecatedRawLocalFileStatus{path=file:/root/en002/Romeo.txt; isDirectory=false; length=145397; replication=1; blocksize=33554432; modification_time=1495535613000; access_time=0; owner=; group=; permission=rw-rw-rw-; isSymlink=false}
inputFile.path.name: shengjing.txt

inputFile.path.name: at.txt
*******
inputFile.path.name: abc.txt

inputFile.path.name: a.txt

inputFile.path.name: av.txt
*******
inputFile.path.name: David.txt
******************************************************************
inputFile.path.name: Oliver.txt
***************************
inputFile.path.name: Jane.txt
**********************************
inputFile.path.name: Romeo.txt
**
**

可能出现的问题:

  1. LocalFileSystem.listStatus返回为 empty, 可能的原因有:
    • Hadoop没有权限读取目录或目录中的文件
    • 本地目录中的文件是中文的,而系统不支持显示中文

FileSystem

DOCS-FileSystem
FileSystem 是Hadoop提供的操作本地文件和HDFS中文件的API,可以实现CRUD等操作

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
/**
* Append to an existing file.
*/
public FSDataOutputStream append(Path f) throws IOException
/**
* Concat existing files together.
*/
public void concat(Path trg,Path[] psrcs)
throws IOException

/**
* Renames Path src to Path dst.
* Can take place on local fs or remote DFS.
*/
public abstract boolean rename(Path src, Path dst) throws IOException




创建目录

1
2
3
4
5
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop/data/20130709");
fs.create(path);
fs.close();

删除目录

1
2
3
4
5
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop/data/20130710");
fs.delete(path);
fs.close();

写文件

1
2
3
4
5
6
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop/data/write.txt");
FSDataOutputStream out = fs.create(path);
out.writeUTF("da jia hao,cai shi zhen de hao!");
fs.close();

读文件

1
2
3
4
5
6
7
8
9
10
11
12
13
14
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop/data/write.txt");

if(fs.exists(path)){
FSDataInputStream is = fs.open(path);
FileStatus status = fs.getFileStatus(path);
// status.getLen() 返回文件的长度 字节长度
byte[] buffer = new byte[Integer.parseInt(String.valueOf(status.getLen()))];
is.readFully(0, buffer);
is.close();
fs.close();
System.out.println(buffer.toString());
}

上传本地文件到HDFS

1
2
3
4
5
6
7
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
Path src = new Path("/home/hadoop/word.txt");
Path dst = new Path("/user/hadoop/data/");
// FileSystem可以直接将本地文件复制到HDFS的制定目录
fs.copyFromLocalFile(src, dst);
fs.close();

删除文件

1
2
3
4
5
6
Configuration conf = new Configuration();  
FileSystem fs = FileSystem.get(conf);
// 删除文件和目录均是通过FileSystem.delete(Path path)
Path path = new Path("/user/hadoop/data/word.txt");
fs.delete(path);
fs.close();

获取给定目录下的所有子目录以及子文件

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
Configuration conf = new Configuration(); 
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop");
getFile(path,fs);
//fs.close();

...

public static void getFile(Path path,FileSystem fs) throws IOException {
FileStatus[] fileStatus = fs.listStatus(path);
for(int i=0;i<fileStatus.length;i++){
if(fileStatus[i].isDir()){
Path p = new Path(fileStatus[i].getPath().toString());
getFile(p,fs);
}else{
System.out.println(fileStatus[i].getPath().toString());
}
}
}

查找某个文件在HDFS集群的位置

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
/** 
* 查找某个文件在HDFS集群的位置
*/
public static void getFileLocal() throws IOException{
Configuration conf = new Configuration();
FileSystem fs = FileSystem.get(conf);
Path path = new Path("/user/hadoop/data/write.txt");

FileStatus status = fs.getFileStatus(path);
BlockLocation[] locations = fs.getFileBlockLocations(status, 0, status.getLen());

int length = locations.length;
for(int i=0;i<length;i++){
String[] hosts = locations[i].getHosts();
System.out.println("block_" + i + "_location:" + hosts[i]);
}
}

HDFS集群上所有节点名称信息

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
/** 
* HDFS集群上所有节点名称信息
*/
public static void getHDFSNode() throws IOException{
Configuration conf = new Configuration();
FileSystem fs = FileSystem.get(conf);

DistributedFileSystem dfs = (DistributedFileSystem)fs;
DatanodeInfo[] dataNodeStats = dfs.getDataNodeStats();

for(int i=0;i<dataNodeStats.length;i++){
System.out.println("DataNode_" + i + "_Node:" + dataNodeStats[i].getHostName());
}

}

读写

FSDataInputStreamFSDataOutputStream

FSDataInputStream 扩展了 DataInputStream 以支持随机读

InputFormat

  • TextInputFormat 是 InputFormat 的默认实现, TextInputFormat
    返回的键是每行的字节偏移量,返回的值是该行的数据。
    key: LongWritable
    value: Text

  • KeyValueTextInputFormat 使用分隔符分割每行,分隔符之前的是键,之后的是值。默认的分隔符是制表符(\T),分离器的属性通过
    key.value.separator.in.input.line中指定

key: Text
value: Text

  • SequenceFileInputFormat<K,V> 用户自定义的序列化格式, 序列化文件为Hadoop专用的压缩二进制文件格式

key: K 用户定义
value: V 用户定义

  • NLineInputFormat key为分片的偏移量,value为包含N行数据的片段,N通过属性 mapred.line.inout.format.linespermap中指定,默认为1
    key: LongWritable
    value: Text

【参考文献】

[<<] Hadoop专题

上一篇中介绍了在Docker搭建三个节点的Hadoop集群,本文主要介绍hadoop集群节点的配置

NameNode
SecondNameNode
DataNode

JobTracker
JobTask

移动计算 而非 移动数据

0.1 查看classpath

hadoop命令行

1
2
hadoop classpath
/opt/module/hadoop/etc/hadoop:/opt/module/hadoop/share/hadoop/common/lib/*:/opt/module/hadoop/share/hadoop/common/*:/opt/module/hadoop/share/hadoop/hdfs:/opt/module/hadoop/share/hadoop/hdfs/lib/*:/opt/module/hadoop/share/hadoop/hdfs/*:/opt/module/hadoop/share/hadoop/mapreduce/lib/*:/opt/module/hadoop/share/hadoop/mapreduce/*:/opt/module/hadoop/share/hadoop/yarn:/opt/module/hadoop/share/hadoop/yarn/lib/*:/opt/module/hadoop/share/hadoop/yarn/*

配置文件

1
2
3
4
<property>
<name>yarn.application.classpath</name>
<value>/opt/module/hadoop/etc/hadoop:/opt/module/hadoop/share/hadoop/common/lib/*:/opt/module/hadoop/share/hadoop/common/*:/opt/module/hadoop/share/hadoop/hdfs:/opt/module/hadoop/share/hadoop/hdfs/lib/*:/opt/module/hadoop/share/hadoop/hdfs/*:/opt/module/hadoop/share/hadoop/mapreduce/lib/*:/opt/module/hadoop/share/hadoop/mapreduce/*:/opt/module/hadoop/share/hadoop/yarn:/opt/module/hadoop/share/hadoop/yarn/lib/*:/opt/module/hadoop/share/hadoop/yarn/*</value>
</property>

[参考文献]

1.

从本篇开始,陆续对常用大数据平台的原理和部分核心源码进行解读和整理。

这里的常用大数据平台包括Hadoop(HDFS、MR)、Spark、Kylin、Flink、HBase、Flume、Elastic Search等

Hadoop专题

  1. 在Docker上搭建Hadoop集群container
  2. Hadoop集群节点配置及搭建
  3. HDFS Java API

HDFS架构

  • HDFS的架构
  • 数据存储与交互

MR架构

对于MapReduce作业,完整的作业运行流程,这里借用刘军老师的Hadoop大数据处理中的一张图:

hadoop

完整过程应该是分为7部分,分别是:

  1. 作业启动:开发者通过控制台启动作业;
  2. 作业初始化:这里主要是切分数据、创建作业和提交作业,与第三步紧密相联;
  3. 作业/任务调度:对于1.0版的Hadoop来说就是JobTracker来负责任务调度,对于2.0版的Hadoop来说就是Yarn中的Resource Manager负责整个系统的资源管理与分配,
    Yarn可以参考IBM的一篇博客Hadoop新MapReduce框架Yarn详解
  4. Map任务;
  5. Shuffle;
  6. Reduce任务;
  7. 作业完成:通知开发者任务完成。

而这其中最主要的MapReduce过程,主要是第4、5、6步三部分,这也是本篇博客重点讨论的地方,详细作用如下:

  1. Map: 数据输入,做初步的处理,输出形式的中间结果;
  2. Shuffle: 按照partition、key对中间结果进行排序合并,输出给reduce线程;
  3. Reduce: 对相同key的输入进行最终的处理,并将结果写入到文件中。

这里先给出官网上关于这个过程的经典流程图:

mapreduce
mapreduce

上图是把MapReduce过程分为两个部分,而实际上从两边的Map和Reduce到中间的那一大块都属于Shuffle过程,也就是说,Shuffle过程有一部分是在Map端,有一部分是在Reduce端,下文也将会分两部分来介绍Shuffle过程。

大部分的情况下,map task与reduce task的执行是分布在不同的节点上的,因此,很多情况下,reduce执行时需要跨节点去拉取其他节点上的map task结果,这样造成了集群内部的网络资源消耗很严重,而且在节点的内部,相比于内存,磁盘IO对性能的影响是非常严重的。如果集群中运行的作业有很多,那么task的执行对于集群内部网络的资源消费非常大。因此,我们对于MapRedue作业Shuffle过程的期望是:

  • 完整地从map task端拉取数据到Reduce端;
  • 在跨节点拉取数据时,尽可能地减少对带宽的不必要消耗;
  • 减少磁盘IO对task执行的影响。

MapReduce之Shuffle过程详述

Map

在进行海量数据处理时,外存文件数据I/O访问会成为一个制约系统性能的瓶颈,因此,Hadoop的Map过程实现的一个重要原则就是:近数据计算,这里主要指两个方面:

  1. 代码靠近数据:
    • 原则:本地化数据处理(locality),即一个计算节点尽可能处理本地磁盘上所存储的数据;
    • 尽量选择数据所在DataNode启动Map任务;
    • 这样可以减少数据通信,提高计算效率;
  2. 数据靠近代码:
    • 当本地没有数据处理时,尽可能从同一机架或最近其他节点传输数据进行处理(host选择算法)。

下面,我们分块去介绍Hadoop的Map过程,map的经典流程图如下:

map-shuffle

输入

  1. map task只读取split分片,split与block(hdfs的最小存储单位,默认为64MB)可能是一对一也能是一对多,但是对于一个split只会对应一个文件的一个block或多个block,不允许一个split对应多个文件的多个block;
  2. 这里切分和输入数据时会涉及到InputFormat的文件切分算法和host选择算法。

文件切分算法

文件切分算法,主要用于确定InputSplit的个数以及每个InputSplit对应的数据段。FileInputFormat以文件为单位切分生成InputSplit,对于每个文件,由以下三个属性值决定其对应的InputSplit的个数:、、、🤣、、、、、😍😱😱😮

  • goalSize: 它是根据用户期望的InputSplit数目计算出来的,即totalSize/numSplits。其中,totalSize为文件的总大小;numSplits为用户设定的Map Task个数,默认情况下是1;
  • minSize:InputSplit的最小值,由配置参数mapred.min.split.size确定,默认是1;
  • blockSize:文件在hdfs中存储的block大小,不同文件可能不同,默认是64MB。

这三个参数共同决定InputSplit的最终大小,计算方法如下:

1
splitSize=max{minSize, min{gogalSize,blockSize}}

host选择算法

FileInputFormat的host选择算法参考《Hadoop技术内幕-深入解析MapReduce架构设计与实现原理》的p50.

Partition

  • 作用:将map的结果发送到相应的reduce端,总的partition的数目等于reducer的数量。
  • 实现功能:
    1. map输出的是key/value对,决定于当前的mapper的part交给哪个reduce的方法是:
      mapreduce提供的Partitioner接口,对key进行hash后,再以reducetask数量取模,然后到指定的job上(HashPartitioner,可以通过job.setPartitionerClass(MyPartition.class)自定义)。
    2. 然后将数据写入到内存缓冲区,缓冲区的作用是批量收集map结果,减少磁盘IO的影响。key/value对以及Partition的结果都会被写入缓冲区。在写入之前,key与value值都会被序列化成字节数组。
  • 要求:负载均衡,效率;

spill(溢写):sort & combiner

  • 作用:把内存缓冲区中的数据写入到本地磁盘,在写入本地磁盘时先按照partition、再按照key进行排序(quick sort);
  • 注意:
    1. 这个spill是由另外单独的线程来完成,不影响往缓冲区写map结果的线程;
    2. 内存缓冲区默认大小限制为100MB,它有个溢写比例(spill.percent),默认为0.8,当缓冲区的数据达到阈值时,溢写线程就会启动,先锁定这80MB的内存,执行溢写过程,map task的输出结果还可以往剩下的20MB内存中写,互不影响。然后再重新利用这块缓冲区,因此Map的内存缓冲区又叫做环形缓冲区(两个指针的方向不会变,下面会详述);
    3. 在将数据写入磁盘之前,先要对要写入磁盘的数据进行一次排序操作,先按<key,value,partition>中的partition分区号排序,然后再按key排序,这个就是sort操作,最后溢出的小文件是分区的,且同一个分区内是保证key有序的;

combine:执行combine操作要求开发者必须在程序中设置了combine(程序中通过job.setCombinerClass(myCombine.class)自定义combine操作)。

  • 程序中有两个阶段可能会执行combine操作:
    1. map输出数据根据分区排序完成后,在写入文件之前会执行一次combine操作(前提是作业中设置了这个操作);
    2. 如果map输出比较大,溢出文件个数大于3(此值可以通过属性min.num.spills.for.combine配置)时,在merge的过程(多个spill文件合并为一个大文件)中还会执行combine操作;
  • combine主要是把形如<aa,1>,<aa,2>这样的key值相同的数据进行计算,计算规则与reduce一致,比如:当前计算是求key对应的值求和,则combine操作后得到<aa,3>这样的结果。
  • 注意事项:不是每种作业都可以做combine操作的,只有满足以下条件才可以:
    1. reduce的输入输出类型都一样,因为combine本质上就是reduce操作;
    2. 计算逻辑上,combine操作后不会影响计算结果,像求和就不会影响;

merge

  • merge过程:当map很大时,每次溢写会产生一个spill_file,这样会有多个spill_file,而最终的一个map task输出只有一个文件,因此,最终的结果输出之前会对多个中间过程进行多次溢写文件(spill_file)的合并,此过程就是merge过程。也即是,待Map Task任务的所有数据都处理完后,会对任务产生的所有中间数据文件做一次合并操作,以确保一个Map Task最终只生成一个中间数据文件。
  • 注意:
    1. 如果生成的文件太多,可能会执行多次合并,每次最多能合并的文件数默认为10,可以通过属性min.num.spills.for.combine配置;
    2. 多个溢出文件合并时,会进行一次排序,排序算法是多路归并排序
    3. 是否还需要做combine操作,一是看是否设置了combine,二是看溢出的文件数是否大于等于3;
    4. 最终生成的文件格式与单个溢出文件一致,也是按分区顺序存储,并且输出文件会有一个对应的索引文件,记录每个分区数据的起始位置,长度以及压缩长度,这个索引文件名叫做file.out.index

内存缓冲区

  1. 在Map Task任务的业务处理方法map()中,最后一步通过OutputCollector.collect(key,value)context.write(key,value)输出Map Task的中间处理结果,在相关的collect(key,value)方法中,会调用Partitioner.getPartition(K2 key, V2 value, int numPartitions)方法获得输出的key/value对应的分区号(分区号可以认为对应着一个要执行Reduce Task的节点),然后将<key,value,partition>暂时保存在内存中的MapOutputBuffe内部的环形数据缓冲区,该缓冲区的默认大小是100MB,可以通过参数io.sort.mb来调整其大小。
  2. 当缓冲区中的数据使用率达到一定阀值后,触发一次Spill操作,将环形缓冲区中的部分数据写到磁盘上,生成一个临时的Linux本地数据的spill文件;然后在缓冲区的使用率再次达到阀值后,再次生成一个spill文件。直到数据处理完毕,在磁盘上会生成很多的临时文件。
  3. 缓存有一个阀值比例配置,当达到整个缓存的这个比例时,会触发spill操作;触发时,map输出还会接着往剩下的空间写入,但是写满的空间会被锁定,数据溢出写入磁盘。当这部分溢出的数据写完后,空出的内存空间可以接着被使用,形成像环一样的被循环使用的效果,所以又叫做环形内存缓冲区
  4. MapOutputBuffe内部存数的数据采用了两个索引结构,涉及三个环形内存缓冲区。下来看一下两级索引结构:

buffer

写入到缓冲区的数据采取了压缩算法
这三个环形缓冲区的含义分别如下:

  1. kvoffsets缓冲区:也叫偏移量索引数组,用于保存key/value信息在位置索引 kvindices 中的偏移量。当 kvoffsets 的使用率超过 io.sort.spill.percent (默认为80%)后,便会触发一次 SpillThread 线程的“溢写”操作,也就是开始一次 Spill 阶段的操作。
  2. kvindices缓冲区:也叫位置索引数组,用于保存 key/value 在数据缓冲区 kvbuffer 中的起始位置。
  3. kvbuffer即数据缓冲区:用于保存实际的 key/value 的值。默认情况下该缓冲区最多可以使用 io.sort.mb 的95%,当 kvbuffer 使用率超过 io.sort.spill.percent (默认为80%)后,便会出发一次 SpillThread 线程的“溢写”操作,也就是开始一次 Spill 阶段的操作。

写入到本地磁盘时,对数据进行排序,实际上是对kvoffsets这个偏移量索引数组进行排序。

Reduce

Reduce过程的经典流程图如下:

reduce-shuffle

copy过程

  • 作用:拉取数据;
  • 过程:Reduce进程启动一些数据copy线程(Fetcher),通过HTTP方式请求map task所在的TaskTracker获取map task的输出文件。因为这时map task早已结束,这些文件就归TaskTracker管理在本地磁盘中。
  • 默认情况下,当整个MapReduce作业的所有已执行完成的Map Task任务数超过Map Task总数的5%后,JobTracker便会开始调度执行Reduce Task任务。然后Reduce Task任务默认启动mapred.reduce.parallel.copies(默认为5)个MapOutputCopier线程到已完成的Map Task任务节点上分别copy一份属于自己的数据。 这些copy的数据会首先保存的内存缓冲区中,当内冲缓冲区的使用率达到一定阀值后,则写到磁盘上。

内存缓冲区

  • 这个内存缓冲区大小的控制就不像map那样可以通过io.sort.mb来设定了,而是通过另外一个参数来设置:mapred.job.shuffle.input.buffer.percent(default 0.7), 这个参数其实是一个百分比,意思是说,shuffile在reduce内存中的数据最多使用内存量为:0.7 × maxHeap of reduce task
  • 如果该reduce task的最大heap使用量(通常通过mapred.child.java.opts来设置,比如设置为-Xmx1024m)的一定比例用来缓存数据。默认情况下,reduce会使用其heapsize的70%来在内存中缓存数据。如果reduce的heap由于业务原因调整的比较大,相应的缓存大小也会变大,这也是为什么reduce用来做缓存的参数是一个百分比,而不是一个固定的值了。

merge过程

  • Copy过来的数据会先放入内存缓冲区中,这里的缓冲区大小要比 map 端的更为灵活,它基于 JVM 的heap size设置,因为 Shuffle 阶段 Reducer 不运行,所以应该把绝大部分的内存都给 Shuffle 用。
  • 这里需要强调的是,merge 有三种形式:1)内存到内存 2)内存到磁盘 3)磁盘到磁盘。默认情况下第一种形式是不启用的。当内存中的数据量到达一定阈值,就启动内存到磁盘的 merge(图中的第一个merge,之所以进行merge是因为reduce端在从多个map端copy数据的时候,并没有进行sort,只是把它们加载到内存,当达到阈值写入磁盘时,需要进行merge) 。这和map端的很类似,这实际上就是溢写的过程,在这个过程中如果你设置有Combiner,它也是会启用的,然后在磁盘中生成了众多的溢写文件,这种merge方式一直在运行,直到没有 map 端的数据时才结束,然后才会启动第三种磁盘到磁盘的 merge (图中的第二个merge)方式生成最终的那个文件。
  • 在远程copy数据的同时,Reduce Task在后台启动了两个后台线程对内存和磁盘上的数据文件做合并操作,以防止内存使用过多或磁盘生的文件过多。

reducer的输入文件

  • merge的最后会生成一个文件,大多数情况下存在于磁盘中,但是需要将其放入内存中。当reducer 输入文件已定,整个 Shuffle 阶段才算结束。然后就是 Reducer 执行,把结果放到 HDFS 上。

shuffle

maptask并行度决定机制

maptask 的并行度决定 map 阶段的任务处理并发度,进而影响到整个 job 的处理速度 那么, mapTask 并行实例是否越多越好呢?其并行度又是如何决定呢?

一个 job 的 map 阶段并行度由客户端在提交 job 时决定, 客户端对 map 阶段并行度的规划
的基本逻辑为:
将待处理数据执行逻辑切片(即按照一个特定切片大小,将待处理数据划分成逻辑上的多 个 split),然后每一个 split 分配一个 mapTask 并行实例处理
这段逻辑及形成的切片规划描述文件,是由 FileInputFormat实现类的 getSplits()方法完成的。
该方法返回的是 List, InputSplit 封装了每一个逻辑切片的信息,包括长度和位置 信息,而 getSplits()方法返回一组 InputSplit

4、切片机制

5、maptask并行度经验之谈

如果硬件配置为 2*12core + 64G,恰当的 map 并行度是大约每个节点 20-100 个 map,最好 每个 map 的执行时间至少一分钟。
(1)如果 job 的每个 map 或者 reduce task 的运行时间都只有 30-40 秒钟,那么就减少该 job 的 map 或者 reduce 数,每一个 task(map|reduce)的 setup 和加入到调度器中进行调度,这个 中间的过程可能都要花费几秒钟,所以如果每个 task 都非常快就跑完了,就会在 task 的开
始和结束的时候浪费太多的时间。
配置 task 的 JVM 重用可以改善该问题:
mapred.job.reuse.jvm.num.tasks,默认是 1,表示一个 JVM 上最多可以顺序执行的 task 数目(属于同一个 Job)是 1。也就是说一个 task 启一个 JVM。这个值可以在 mapred-site.xml 中进行更改, 当设置成多个,就意味着这多个 task 运行在同一个 JVM 上,但不是同时执行,
是排队顺序执行
(2)如果 input 的文件非常的大,比如 1TB,可以考虑将 hdfs 上的每个 blocksize 设大,比如 设成 256MB 或者 512MB
6、reducetask并行度决定机制

文件切分算法,主要用于确定InputSplit的个数以及每个InputSplit对应的数据段。FileInputFormat以文件为单位切分生成InputSplit,对于每个文件,由以下三个属性值决定其对应的InputSplit的个数:

goalSize: 它是根据用户期望的InputSplit数目计算出来的,即totalSize/numSplits。其中,totalSize为文件的总大小;numSplits为用户设定的Map Task个数,默认情况下是1;
minSize:InputSplit的最小值,由配置参数mapred.min.split.size确定,默认是1;
blockSize:文件在hdfs中存储的block大小,不同文件可能不同,默认是64MB。
这三个参数共同决定InputSplit的最终大小,计算方法如下:

splitSize=max{minSize, min{gogalSize,blockSize}}

1
2
3
4
1. 用户设置了numSplit,那么goalSize=totalSize/numSplit
2. minSize=max(1,minSplitSize)
3. splitSize=max(minSplitSize, min(goalSize,blockSize))
4. task个数=totalSize除以splitSize

FileInputFormat的host选择算法参考《Hadoop技术内幕-深入解析MapReduce架构设计与实现原理》的p50.

yarn

Hadoop MapReduce 2.x 工作原理
Hadoop新MapReduce框架Yarn详解

hadoop 优缺点

与许多伟大的系统一样,Hadoop开启了我们对一个新的空间的问题。具体来说,Hadoop在存储和提供对大量数据的访问方面表现优异,然而,它并没有提供关于数据访问速度的性能保证。此外,尽管Hadoop是高可用性系统,但在高并发负载下性能下降。最后,虽然Hadoop在存储数据方面表现良好,但它并未针对提取数据和使数据立即可读而进行优化。