本文主要分析了hadoop客戶端read和write block的流程. 以及client和datanode通信的協(xié)議, 數(shù)據(jù)流格式等.
hadoop客戶端與namenode通信通過RPC協(xié)議, 但是client 與datanode通信并沒有使用RPC, 而是直接使用socket, 其中讀寫時的協(xié)議也不同, 本文分析了hadoop 0.20.2版本的(0.19版本也是一樣的)client與datanode通信的原理與通信協(xié)議. 另外需要強(qiáng)調(diào)的是0.23及以后的版本中client與datanode的通信協(xié)議有所變化, 使用了protobuf作為序列化方式.
Write block
1. 客戶端首先通過namenode.create, 向namenode請求創(chuàng)建文件, 然后啟動dataStreamer線程
2. client包括三個線程, main線程負(fù)責(zé)把本地數(shù)據(jù)讀入內(nèi)存, 并封裝為Package對象, 放到隊列dataQueue中.
3. dataStreamer線程檢測隊列dataQueue是否有package, 如果有, 則先創(chuàng)建BlockOutPutStream對象(一個block創(chuàng)建一次, 一個block可能包括多個package), 創(chuàng)建的時候會和相應(yīng)的datanode通信, 發(fā)送DATA_TRANSFER_HEADER信息并獲取返回. 然后創(chuàng)建ResponseProcessor線程, 負(fù)責(zé)接收datanode的返回ack確認(rèn)信息, 并進(jìn)行錯誤處理.
4. dataStreamer從dataQueue中拿出Package對象, 發(fā)送給datanode. 然后繼續(xù)循環(huán)判斷dataQueue是否有數(shù)據(jù)…..
下圖展示了write block的流程.
下圖是報文的格式
Read block
主要在BlockReader類中實現(xiàn).
初始化newBlockReader時,
1. 通過傳入?yún)?shù)sock創(chuàng)建new SocketOutputStream(socket, timeout), 然后寫通信信息, 與寫block的header不大一樣.
//write the header.
out.writeShort( DataTransferProtocol.DATA_TRANSFER_VERSION );
out.write( DataTransferProtocol.OP_READ_BLOCK );
out.writeLong( blockId );
out.writeLong( genStamp );
out.writeLong( startOffset );
out.writeLong( len );
Text.writeString(out, clientName);
out.flush();
2. 創(chuàng)建輸入流 new SocketInputStream(socket, timeout)
3. 判斷返回消息 in.readShort() != DataTransferProtocol.OP_STATUS_SUCCESS
4. 根據(jù)輸入流創(chuàng)建checksum : DataChecksum checksum = DataChecksum.newDataChecksum( in )
5. 讀取第一個Chunk的位置: long firstChunkOffset = in.readLong()
注: 512個字節(jié)為一個chunk計算checksum(4個字節(jié))
6. 接下來在BlockReader的read方法中讀取具體數(shù)據(jù): result = readBuffer(buf, off, realLen)
7. 一個一個chunk的讀取
int packetLen = in.readInt();
long offsetInBlock = in.readLong();
long seqno = in.readLong();
boolean lastPacketInBlock = in.readBoolean();
int dataLen = in.readInt();
IOUtils.readFully(in, checksumBytes.array(), 0,
checksumBytes.limit());
IOUtils.readFully(in, buf, offset, chunkLen);
8. 讀取數(shù)據(jù)后checksum驗證; FSInputChecker.verifySum(chunkPos)
新聞熱點
疑難解答
圖片精選