2017年10月18日星期三

其它技术备忘

G1 gc

1.garbage first是指内存分很多小区域,每个区域可以标记垃圾比例,回收的时候找比例高的回收
2.为啥能控制停顿时间?能自动调整参数是一个原因,还有一个原因是能控制回收分区的数量
3.为啥能解决碎片?同样是因为分区,回收的时候是把一些分区挪移到其它的分区,同时做了压缩。核心创新还是分区

java性能权威指南的几个点

1.TLAB线程本地对象分配区适当调整可能会优化对象分配速度
2.linux huge使用加快page寻址,避免page换出
3.默认开启指针压缩但最大只能支持到32g,一般内存不要超过这个值
4.java8 parallelStream 会做自动并行化,使用恰当可以优化性能
5.@Contended可以对齐cache line,而不用自己去推算cpu cache line大小

关于mmap的性能

1.mmap快在不用每次读写都拷贝内存和系统调用
2.写的时候尤其是顺序写,如果每次写入的大小比较大,内存拷贝和系统调用的优势在实际io面前就没那么大了,所以mmap和fwrite的性能实际差距不太大
3.一般的建议是append的方式写文件用fwirte,因为mmap resize需要反复撤销和创建mmap这个成本较高没必要。不过如果一开始就知道文件大小而且很大的情况下用mmap也可以。如果读比较多就比较适合用mmap。比较经典的例子就是kafka和metaq,kafka只有索引用了mmap而metaq是索引和log文件都用了mmap。可以看看linus大神怎么说https://stackoverflow.com/questions/35891525/mmap-for-writing-sequential-log-file-for-speed

io_uring

1.简单描述就是有submission queue和completion queue两个队列是用户态和内核态共享,所以用户态提交io请求由内核线程负责处理sq并返回cq,整个过程没有系统调用


2017年10月9日星期一

kafka精确一次投递和事务消息学习整理

概述

今年6月发布的kafka 0.11.0.0包含两个比较大的特性,exactly once delivery和transactional transactional messaging。之前一直对事务这块比较感兴趣,所以抽空详细学习了一下,感觉收获还是挺多的。

对这两个特性的详细描述可以看这三篇文档,
https://cwiki.apache.org/confluence/display/KAFKA/KIP-98+-+Exactly+Once+Delivery+and+Transactional+Messaging
https://cwiki.apache.org/confluence/display/KAFKA/Idempotent+Producer
https://cwiki.apache.org/confluence/display/KAFKA/Transactional+Messaging+in+Kafka

需求场景


精确一次投递

消息重复一直是消息领域的一个痛点,而消息重复可能发生于下面这些场景
1.消息发送端发出消息,服务端落盘以后因为网络等种种原因发送端得到一个发送失败的响应,然后发送端重发消息导致消息重复。  
2.消息消费端在消费过程中挂掉另一个消费端启动拿之前记录的位点开始消费,由于位点的滞后性可能会导致新启动的客户端有少量重复消费。

先说说问题2,一般的解决方案是让下游做幂等或者尽量每消费一条消息都记位点,对于少数严格的场景可能需要把位点和下游状态更新放在同一个数据库里面做事务来保证精确的一次更新或者在下游数据表里面同时记录消费位点,然后更新下游数据的时候用消费位点做乐观锁拒绝掉旧位点的数据更新。

问题1的解决方案也就是kafka实现的方案是每个producer有一个producer id,服务端会通过这个id关联记录每个producer的状态,每个producer的每条消息会带上一个递增的sequence,服务端会记录每个producer对应的当前最大sequence,如果新的消息带上的sequence不大于当前的最大sequence就拒绝这条消息,问题1的场景如果消息落盘会同时更新最大sequence,这个时候重发的消息会被服务端拒掉从而避免消息重复。后面展开详细说一下这个解决方案。

事务消息

1.最简单的需求是producer发的多条消息组成一个事务这些消息需要对consumer同时可见或者同时不可见
2.producer可能会给多个topic,多个partition发消息,这些消息也需要能放在一个事务里面,这就形成了一个典型的分布式事务
3.kafka的应用场景经常是应用先消费一个topic,然后做处理再发到另一个topic,这个consume-transform-produce过程需要放到一个事务里面,比如在消息处理或者发送的过程中如果失败了,消费位点也不能提交
4.producer或者producer所在的应用可能会挂掉,新的producer启动以后需要知道怎么处理之前未完成的事务
5.流式处理的拓扑可能会比较深,如果下游只有等上游消息事务提交以后才能读到,可能会导致rt非常长吞吐量也随之下降很多,所以需要实现read committed和read uncommitted两种事务隔离级别

一个比较典型的consume-transform-produce的场景像下面这样
public class KafkaTransactionsExample {
  
  public static void main(String args[]) {
    KafkaConsumer consumer = new KafkaConsumer<>(consumerConfig);
    KafkaProducer producer = new KafkaProducer<>(producerConfig);

    producer.initTransactions();
     
    while(true) {
      ConsumerRecords records = consumer.poll(CONSUMER_POLL_TIMEOUT);
      if (!records.isEmpty()) {
        producer.beginTransaction();
        List> outputRecords = processRecords(records);
        for (ProducerRecord outputRecord : outputRecords) {
          producer.send(outputRecord);
        }
        sendOffsetsResult = producer.sendOffsetsToTransaction(getUncommittedOffsets());
        producer.endTransaction();
      }
    }
  }
}

几个关键概念和推导

1.因为producer发送消息可能是分布式事务,所以引入了常用的2PC,所以有事务协调者(Transaction Coordinator)。Transaction Coordinator和之前为了解决脑裂和惊群问题引入的Group Coordinator在选举和failover上面类似。
2.事务管理中事务日志是必不可少的,kafka使用一个内部topic来保存事务日志,这个设计和之前使用内部topic保存位点的设计保持一致。事务日志是Transaction Coordinator管理的状态的持久化,因为不需要回溯事务的历史状态,所以事务日志只用保存最近的事务状态。
3.因为事务存在commit和abort两种操作,而客户端又有read committed和read uncommitted两种隔离级别,所以消息队列必须能标识事务状态,这个被称作Control Message。
4.producer挂掉重启或者漂移到其它机器需要能关联的之前的未完成事务所以需要有一个唯一标识符来进行关联,这个就是TransactionalId,一个producer挂了,另一个有相同TransactionalId的producer能够接着处理这个事务未完成的状态。注意不要把TransactionalId和数据库事务中常见的transaction id搞混了,kafka目前没有引入全局序,所以也没有transaction id,这个TransactionalId是用户提前配置的。
5. TransactionalId能关联producer,也需要避免两个使用相同TransactionalId的producer同时存在,所以引入了producer epoch来保证对应一个TransactionalId只有一个活跃的producer epoch

重要的类图




架构和数据流

事务消息

官方文档的数据流组件图

上图中每个方框代表一台独立的机器,图下方比较长的圆角矩形代表kafka topic,图中间的两个角是圆的的方框代表broker里面的逻辑组件,箭头代表rpc调用。

接下来说一下事务的数据流,这里基本按照官方文档的结构加上我自己看代码的一点补充

1.首先producer需要找到transaction coordinator。TransactionManager.lookupCoordinator
private synchronized void lookupCoordinator(FindCoordinatorRequest.CoordinatorType type, String coordinatorKey) {
    switch (type) {
        case GROUP:
            consumerGroupCoordinator = null;
            break;
        case TRANSACTION:
            transactionCoordinator = null;
            break;
        default:
            throw new IllegalStateException("Invalid coordinator type: " + type);
    }

    FindCoordinatorRequest.Builder builder = new FindCoordinatorRequest.Builder(type, coordinatorKey);
    enqueueRequest(new FindCoordinatorHandler(builder));
}

2.获取producer id,producer id是比较重要的概念,精确一次投递需要producer id+sequence防止重复投递,事务消息也需要保存transactional id和producer id的对应关系。客户端调用KafkaProducer.initTransactions的时候会向coordinator请求producer id,TransactionManager.initializeTransactions
public synchronized TransactionalRequestResult initializeTransactions() {
    ensureTransactional();
    transitionTo(State.INITIALIZING);
    setProducerIdAndEpoch(ProducerIdAndEpoch.NONE);
    this.sequenceNumbers.clear();
    InitProducerIdRequest.Builder builder = new InitProducerIdRequest.Builder(transactionalId, transactionTimeoutMs);
    InitProducerIdHandler handler = new InitProducerIdHandler(builder);
    enqueueRequest(handler);
    return handler.result;
}

coordinator端处理逻辑在TransactionCoordinator.handleInitProducerId流程比较复杂,首先如果对应的transactional id没有产生过producer id会找producerIdManager生成一个
val coordinatorEpochAndMetadata = txnManager.getTransactionState(transactionalId).right.flatMap {
  case None =>
    val producerId = producerIdManager.generateProducerId()
    val createdMetadata = new TransactionMetadata(transactionalId = transactionalId,
      producerId = producerId,
      producerEpoch = RecordBatch.NO_PRODUCER_EPOCH,
      txnTimeoutMs = transactionTimeoutMs,
      state = Empty,
      topicPartitions = collection.mutable.Set.empty[TopicPartition],
      txnLastUpdateTimestamp = time.milliseconds())
    txnManager.putTransactionStateIfNotExists(transactionalId, createdMetadata)

  case Some(epochAndTxnMetadata) => Right(epochAndTxnMetadata)
}

producer id需要全局唯一,有点类似于tddl sequence的生成逻辑,ProducerIdManager.generateProducerId会一次申请一批id然后在zk上面保存状态,本地每次生成+1,如果超出了当前批次的范围就去找zk重新申请

拿到了producer id接下来处理事务状态,保证之前的事务状态能够处理完毕,该提交的提交,该回滚的回滚。
TransactionCoordinator.prepareInitProduceIdTransit处理producer id的变化比如开始一个新的事务可能会增加producer epoch,也可能生成新的producer id
case PrepareAbort | PrepareCommit =>
  // reply to client and let it backoff and retry
  Left(Errors.CONCURRENT_TRANSACTIONS)

case CompleteAbort | CompleteCommit | Empty =>
  val transitMetadata = if (txnMetadata.isProducerEpochExhausted) {
    val newProducerId = producerIdManager.generateProducerId()
    txnMetadata.prepareProducerIdRotation(newProducerId, transactionTimeoutMs, time.milliseconds())
  } else {
    txnMetadata.prepareIncrementProducerEpoch(transactionTimeoutMs, time.milliseconds())
  }

  Right(coordinatorEpoch, transitMetadata)

case Ongoing =>
  // indicate to abort the current ongoing txn first. Note that this epoch is never returned to the
  // user. We will abort the ongoing transaction and return CONCURRENT_TRANSACTIONS to the client.
  // This forces the client to retry, which will ensure that the epoch is bumped a second time. In
  // particular, if fencing the current producer exhausts the available epochs for the current producerId,
  // then when the client retries, we will generate a new producerId.
  Right(coordinatorEpoch, txnMetadata.prepareFenceProducerEpoch())

最后如果之前的事务处于进行中的状态会回滚事务
handleEndTransaction(transactionalId,
  newMetadata.producerId,
  newMetadata.producerEpoch,
  TransactionResult.ABORT,
  sendRetriableErrorCallback)

或者就是新事务,往事务日志里面插一条日志(对应数据流图中的2a)
txnManager.appendTransactionToLog(transactionalId, coordinatorEpoch, newMetadata, sendPidResponseCallback)

3.客户端调用KafkaProducer.beginTransaction开始新事务。这一步相对简单,就是客户端设置状态成State.IN_TRANSACTION

4.consume-transform-produce过程,这一步是实际消费消息和生产消息的过程。
(1)客户端发送消息时(KafkaProducer.send),对于新碰到的TopicPartition会触发AddPartitionsToTxnRequest。服务端对应的处理在TransactionCoordinator.handleAddPartitionsToTransaction,主要做的事情是更新事务元数据和记录事务日志(对应数据流图中的4.1a)。在事务中记录partition的作用是后面给事务每个partition发送提交或者回滚标记时需要事务所有的partition。
(2)客户端通过KafkaProducer.send发送消息(ProduceRequest),比较早的kafka版本增加了PID,epoch,sequence number等几个字段,对应数据流图中的4.2a
(3)客户端调用KafkaProducer.sendOffsetsToTransaction保存事务消费位点。服务端的处理逻辑在TransactionCoordinator.handleAddPartitionsToTransaction,和4.1基本是一样的,不同的是4.3记录的是记录消费位点的topic(GROUP_METADATA_TOPIC_NAME)。
(4)4.3调用的后半部分会触发TxnOffsetCommitRequest,通过数据消息的方式把消费位点持久化到GROUP_METADATA_TOPIC_NAME(__consumer-offsets)这个topic里面去,对应数据流图中的4.4a。
客户端发起逻辑在AddOffsetsToTxnHandler.handleResponse
if (error == Errors.NONE) {
    log.debug("{}Successfully added partition for consumer group {} to transaction", logPrefix,
            builder.consumerGroupId());

    // note the result is not completed until the TxnOffsetCommit returns
    pendingRequests.add(txnOffsetCommitHandler(result, offsets, builder.consumerGroupId()));
    transactionStarted = true;
}
因为需要处理可见性相关的逻辑,服务端事务消费位点和普通消费位点提交的处理逻辑稍有不同,调用GroupCoordinator.handleTxnCommitOffsets而不是handleCommitOffsets。

5.结束事务需要调用KafkaProducer.commitTransaction或者KafkaProducer.abortTransaction
(1)首先客户端会发送一个EndTxnRequest,而服务端由TransactionCoordinator.handleEndTransaction处理。

handleEndTransaction首先会做一个可能的状态转换让事务进入预提交或者预放弃阶段
else txnMetadata.state match {
  case Ongoing =>
    val nextState = if (txnMarkerResult == TransactionResult.COMMIT)
      PrepareCommit
    else
      PrepareAbort
    Right(coordinatorEpoch, txnMetadata.prepareAbortOrCommit(nextState, time.milliseconds()))

接下来会在事务日志里面记录PREPARE_COMMIT或者PREPARE_ABORT日志,对应数据流图中的5.1a
txnManager.appendTransactionToLog(transactionalId, coordinatorEpoch, newMetadata, sendTxnMarkersCallback)

再接下来会往用户数据日志里面发送COMMIT或者ABORT的Control Message,最后往事务日志里面写入COMMIT或者ABORT,才算完成了事务的提交过程。这个过程是用回调的方式组织起来的,代码的流程是TransactionStateManager.appendTransactionToLog->TransactionMarkerChannelManager.addTxnMarkersToSend->TransactionStateManager.appendTransactionToLog

(2)往用户数据日志里面发送COMMIT或者ABORT的Control Message的过程,对应数据流图中的5.2a
发起点在回调方法sendTxnMarkersCallback,这个方法首先会做转台转换让事务进入CompleteCommit或者CompleteAbort状态
case PrepareCommit =>
  if (txnMarkerResult != TransactionResult.COMMIT)
    logInvalidStateTransitionAndReturnError(transactionalId, txnMetadata.state, txnMarkerResult)
  else
    Right(txnMetadata, txnMetadata.prepareComplete(time.milliseconds()))
case PrepareAbort =>
  if (txnMarkerResult != TransactionResult.ABORT)
    logInvalidStateTransitionAndReturnError(transactionalId, txnMetadata.state, txnMarkerResult)
  else
    Right(txnMetadata, txnMetadata.prepareComplete(time.milliseconds()))

最后会往事务相关的每个broker发送WriteTxnMarkersRequest,如果事务包含消费位点也会往__consumer-offsets所在的broker发请求。 broker端的处理在KafkaApis.handleWriteTxnMarkersRequest会把control message写入日志
replicaManager.appendRecords(
  timeout = config.requestTimeoutMs.toLong,
  requiredAcks = -1,
  internalTopicsAllowed = true,
  isFromClient = false,
  entriesPerPartition = controlRecords,
  responseCallback = maybeSendResponseCallback(producerId, marker.transactionResult))

(3)事务日志写入最终的COMMIT或者ABORT日志,对应数据流图的5.3,这一步完成了一个事务就算彻底完成了。
发起点在回调方法appendToLogCallback

精确一次投递



大致列一下流程中的关键节点

1.在客户端每次发送消息之前会检查是否有producerId如果没有会找服务端去申请,Sender.run
if (transactionManager != null) {
    if (!transactionManager.isTransactional()) {
        // this is an idempotent producer, so make sure we have a producer id
        maybeWaitForProducerId();
    } 
服务端会在TransactionCoordinator.handleInitProducerId处理,前面事务消息提到过

2.在生成消息内容的时候(RecordAccumulator.drain)会获取当前的的sequenceNumber(TransactionManager.sequenceNumber)放到消息体里面。而sequenceNumber的自增是在发送上一批消息返回是触发的(Sender.handleProduceResponse)。

3.broker实际写入消息之前(Log.append)才会对sequenceNumber进行校验,校验的具体逻辑在ProducerStateManager.validateAppend

其他一些技术点

1.sequenceNumber可以设计成每个producer唯一或者更细粒度的对每个topic-partition唯一,topic-partition唯一的好处是对于每个topic-partition sequenceNumber可以设计成连续的,这样broker端可以做更强的校验,比如检查丢消息,kafka使用的就是细粒度的方法,发现sequenceNumber不连续的时候会抛异常OutOfOrderSequenceException

2.发消息(KafkaProducer.doSend)是个异步的过程,但同时提供Future返回值使得在必要的时候可以把异步变成同步等待。kafka也实现了攒消息批量发送的能力(RecordAccumulator.append),攒消息的存放方式是一个大hash map,key是topic-partition,消息实际刷出和发送在一个单独的线程中执行,调用Sender.sendProducerData。被刷出的消息的判定在RecordAccumulator.ready,主要依据是消息集是否已满或者是否超时。

3.Consumer消费的时候怎样控制事务可见性呢?一个比较直观的方法就是先把事务消息buffer起来,然后遇到提交或者回滚标志的时候做相应的处理,kafka处理的更巧妙一些。首先是不能读到未提交的事务的控制,kafka引入了lastStableOffset这个概念,lastStableOffset是当前已经提交的事务的最大位点。在ReplicaManager.readFromLocalLog里面有控制,
val initialHighWatermark = localReplica.highWatermark.messageOffset
val lastStableOffset = if (isolationLevel == IsolationLevel.READ_COMMITTED)
  Some(localReplica.lastStableOffset.messageOffset)
else
  None

// decide whether to only fetch committed data (i.e. messages below high watermark)
val maxOffsetOpt = if (readOnlyCommitted)
  Some(lastStableOffset.getOrElse(initialHighWatermark))
else
  None
...
val fetch = log.read(offset, adjustedFetchSize, maxOffsetOpt, minOneMessage, isolationLevel)

这样未提交的事务对客户端就不可见了

4.还有一个需求是要识别并且跳过那些在aborted事务内的消息,这些消息可能和非事务消息混在一起。kafka读消息的返回信息中会带上本批读取的消息中回滚事务列表来帮助客户端跳过。
case class FetchDataInfo(fetchOffsetMetadata: LogOffsetMetadata,
                         records: Records,
                         firstEntryIncomplete: Boolean = false,
                         abortedTransactions: Option[List[AbortedTransaction]] = None)

回滚事务列表是在读取消息日志(Log.read)的过程中撸的
return isolationLevel match {
  case IsolationLevel.READ_UNCOMMITTED => fetchInfo
  case IsolationLevel.READ_COMMITTED => addAbortedTransactions(startOffset, segmentEntry, fetchInfo)
}
那接下来的问题是kafka是如何快速拿到回滚事务列表的呢?kafka为这件事做了一个文件索引,文件后缀名是'.txnindex',相关的管理逻辑在TransactionIndex。

5.broker端对于producer的状态管理,broker需要记录producer对应的最大sequenceNumber,epoch之类的信息。相关逻辑是放在ProducerStateManager里面的,broker每次写入消息的时候(Log.append)都会更新producer信息(ProducerStateManager.update)。由于只有当有消息写入的时候producer state才会被更新,所以当broker挂掉的时候producer的状态需要被持久化,kafka又弄了一个文件'.snapshot'来持久化producer信息。

设计上的体会

感觉kafka在设计上概念的统一和架构的连贯上做的特别好,比如producerId的引入把精确一次投递和事务消息都给串联起来了。

印象更深刻的例子是kafka早先几个版本依次推出了几个特性,
1.把位点当做普通消息保存
2.加入了消息清理机制,只保留key最新的value
3.加入了broker端的coordinator解决惊群和脑裂问题

然后这几个特性在事务消息这块都用上了,首先把位点当普通消息保存在概念上统一了消息发送和消费,同时消息同步也成了broker之间同步状态的基础机制这样就不用再弄一套状态同步机制了,不过这样做的缺点是只有写消息才能同步broker状态某些特殊情况可能有点小麻烦。利用消息队列保存状态的一个毛病是比较浪费资源,而消息清理机制恰好解决了这个问题。最后是broker端的coordinator机制可以用在consumer group协调者也可以用在事务协调者上面。这种层次递进的特性累加真是相当有美感而且感觉是深思熟虑的结果。

参考资料

https://cwiki.apache.org/confluence/display/KAFKA/KIP-98+-+Exactly+Once+Delivery+and+Transactional+Messaging
https://cwiki.apache.org/confluence/display/KAFKA/Idempotent+Producer
https://cwiki.apache.org/confluence/display/KAFKA/Transactional+Messaging+in+Kafka
https://medium.com/@jaykreps/exactly-once-support-in-apache-kafka-55e1fdd0a35f




2017年1月14日星期六

nodejs源码探险

动机

之所以称之为探险是因为我是一名java程序员,对c和c++的了解仅限于大学时学的一点点知识,而javascript虽然用过一些,但深入程度有限,nodejs的源码组合了c,c++,js对我来说实在难度不小,怎奈我对特别酷的技术没什么抵抗力,所以还是入了坑。而学习nodejs代码的过程也是充满了挫折,断断续续持续了一年的时间而且中途差点就走不下去而放弃了,现在总算有点弄明白了,所以写下这篇文章算是付出的这些艰辛的一点结晶吧。

确立目标

学习一个框架的源代码之前总得对这个框架有所了解,所以先找了本node.js in action大概的翻了一遍,从这本书得到的对node的印象大概是1.node和前端的js没什么区别,但是有个require是前端js没有的 2.node可以写server端程序,可以操作文件,操作网络 3.node是单线程的 4.node的风格是异步,io风格也是异步回调 5.node开发web应用确实很简单方便,而且框架不少。 之后又在网上零碎的了解了一些node知识1.node有三大块--nodejs本身,libuv,v8引擎 2.node的tick, setTimeout行为比较有趣

到这里我学习nodejs的目标就大致确定了
1.学习nodejs本身,看看怎么整合js,v8和libuv
2.学习libuv,了解一下linux下面的网络,io编程,同时也了解一下epoll因为用java nio的时候只能了解到nio api这层就止步了,比较想知道下面发生的事情
3.学习v8引擎,看看这牛逼的引擎为啥这么快,gc和即时编译这些都是怎么做的(事后证明这个目标太贪心了,也导致差点玩下去了)

在上面的目标下要重点了解的技术点也可以列出来
1.nodejs的require是怎么实现的
2.nodejs的js和libuv的c是怎么交互的
3.nextTick和setTimeout这些是怎么实现的
4.libuv的异步io是怎么实现的
5.nodejs里面的domain能捕获相关的异常看着很酷,怎么实现的(这个点还是读源代码过程中发现的)
6.epoll为啥那么牛逼
7.了解c/c++的编程风格

准备工作

玩源代码首先要能debug,所以一顿google把准备工作做足了

调试js端代码

这个相对简单一点,
1.npm install -g node-inspector
2.node --debug-brk /home/huying/development/nodejs/huying-code/basicHttp.js
3.启动node-inspector

调试c和c++代码

这一块就没有官方攻略了,只能自己diy,虽然传说中vim或者emacs+gdb很牛逼,但是就我自己emacs+jdb的体验来文件编辑器做debug操作起来并不舒服,像树形展开这类操作没有鼠标点来方便,所以还是选了比较全能的eclipse

后面的步骤就直接列了
1.去官网下载node源代码,我下的是node-v4.4.7
2.打开源码的debug开关 ./configure --prefix /home/huying/development/nodejs/node-debug-version --gdb 或者直接 ./configure --debug
3.修改Makefile,打开BUILDTYPE=Debug
4.make
5.eclipse首先在源码上创建c/c++项目,然后用node_g创建一个运行项就可以debug了(gdb deubg可以直接gdb --args ./node_g myscript.js)
6.碰到glibc这样看不到源码的,可以先apt-get install源码,然后在source这里用创建目录映射的方式去关联源码

查看v8变量

这里我踩了一个大坑,每次debug进入v8的时候变量被v8包装一下就没法看了,在eclipse里面v8的变量都是Local<String>这样的东东,无论怎么展开都没法看到里面的String长啥样。。。
志在了解v8引擎的我自然不甘心,各种google之下发现可以给gdb开pretty print,可是stl的pretty print好找,v8的pretty print脚本真不多而且看起来比较靠谱的还是用scheme写的。。。这里折腾了不少时间都没看到效果,而且平时只能拿出零碎的时间玩node所以感觉搞不下去了。。。

后来不得已调整了目标,放弃了了解v8,其实了解v8这个目标订的太大即使debug好使,可能陷入其中也会耽误了解node js这个基本目标。但是完全看不到v8里面的变量内容对理解nodejs的某些流程还是会有影响的,后来我找到了一些变通的解决办法,比如c端和js端一起debug,还有就是可以修改源代码,加上这个函数
    static const char* ToCString(Local value) {
     v8::String::Utf8Value string(value);
     char *str = (char *) malloc(string.length() + 1);
     strcpy(str, *string);
     return str;
    }

然后就能打印出v8变量的内容了,虽然使用起来还是不怎么方便


捏个软柿子先

看源代码首先得找个切入点,我一般是从启动流程看起,nodejs的入口在c这段,感觉有点难咧,所以就先看看js这端的启动逻辑吧

1.启动执行的是node.js文件的startup()方法,找到这个入口当然是google来的但是后面看c++部分的时候得到了印证,在node.cc->LoadEnvironment 函数的末尾
 
  Local script_name = FIXED_ONE_BYTE_STRING(env->isolate(), "node.js");
  Local f_value = ExecuteString(env, MainSource(env), script_name);

  Local arg = env->process_object();
  f->Call(global, 1, &arg);

2.上来就是一句var EventEmitter = NativeModule.require('events'),先了解一下NativeModule.require

NativeModule.require->NativeModule.prototype.compile->NativeModule._source = process.binding('natives')->NativeModule.wrap(给source加上一个function(exports, require ...) 这样一个头和尾)->ContextifyScript.runInThisContext(这对应一个c++类在node_contextify.cc)->Local<Script> script = unbound_script->BindToCurrentContext()->script->Run()(这里返回的是一个function,就是wrap加了function头尾的那个)->fn(this.exports, NativeModule.require, this, this.filename)->return nativeModule.exports

这里->只是我思考的顺序不完全代表调用的过程,可以看得出NativeModule.require这个过程还是挺复杂的。require首先通过process.binding从c++那一端拿到module的源代码,然后用wrap方法给加了个头尾,注意这里是个很有趣的元编程技巧,我一直想不通那个长得像关键字的require()函数是从哪里来的,其实就是这个wrap给传进来的,传的值是NativeModule.require
wrap以后的js代码又被传回c++那边给v8引擎执行,执行完以后回到js端返回exports属性

3.接下来一句EventEmitter.call(process)也让我头晕,冷静一点以后发现process其实就是当this给传进去了,然后EventEmitter初始化方法给process对象设置了各种属性

4.接下来各种startup.processXXX都比较好理解,就是设置了各种数据结构,然后是debug相关参数的处理

5.接下来比较重要的一句是startup.preloadModules(),这个是由--require参数触发的预先加载

NativeModule.require('module')._preloadModules(这里比较重要的module.js就登场了)->Module.prototype.require->Module._load->Module._resolveFilename->NativeModule.nonInternalExists->Module.prototype.load->this.paths = Module._nodeModulePaths(path.dirname(filename))(设置查找路径)->Module._extensions[extension](this, filename)(这句看着不起眼,却很重要,用的functional的技巧)->module._compile(internalModule.stripBOM(content), filename)->Module.wrap->runInThisContext(这两个都和NativeModule里面调的一样)->const require = internalModule.makeRequireFunction.call(this)(这句比较诡异,结果还是拿到了module.require函数本身,但是给函数对象加了几个属性)->compiledWrapper.apply(this.exports, args)(和之前NativeModule相同的trick)

整体逻辑就是把参数--require里面的module每个都调一下Module.require,发生的事情也和NativeModule里面的类似。Module和NativeModule的核心逻辑类似,感觉Module的功能更广一些

6.终于走到跑main脚本了,Module.runMain()->Module._load(process.argv[1], null, true) 其实就是load一下目标脚本,后面还跟了一句process._tickCallback,这个是给无孔不入的process.nextTick用的,后面再详细说明

js端启动流程就看得差不多了,这个过程的收获是大致了解了nodejs js端的启动流程,同时也了解了长得像关键字的reqiure的实现,对其中的元编程技巧也是印象深刻。这一段学习过程的方法也比较简单,顺着链路读代码就ok了,对于c++端发生的事情用用grep大致就能分析出来了


初看c++端启动流程

到这里就进入噩梦模式了,刚开始看c++代码的时候最大的问题是把握不好抽象级别,一不小心就扎入一个细节出不来,或者忽略一个关键的细节。第一遍过启动流程的感觉是云里雾里,不过也抓住了几个比较重要的点

1.main(node_main.cc) 启动入口

2.PlatformInit(node.cc)  信号量和文件描述符相关的初始化

3.Init(node.cc) 这里有一个比较重要的调用 uv_default_loop->uv_loop_init(初始化lib_uv的核心数据结构uv_loop_t)。跟踪进去以后看到的是各种QUEUE_INIT(&loop->wq),uv_async_init,uv_signal_init,马上就迷失了,关键问题是不知道这各种宏,各种数据结构都有什么用,所以这里先跳过,等后面经验值够了再回来

4.接下来几句v8::platform::CreateDefaultPlatform,V8::InitializePlatform(default_platform),V8::Initialize()很明显是初始化v8引擎了,查查文档就可以确定是v8的标准用法

5.StartNodeInstance(node.cc) 这里有好多v8相关的东东,所以先得稍微了解一下才能进行下去
只要google一下Isolate,HandleScope 这些关键字就能找到不少讲解这些概念的文章,在这里稍做总结

(1)Isolate代表一个v8 engine实例 ,Context代表一次js执行的上下文
(2)Handle是v8环境的对象句柄,因为指针可能被gc移动,所以必须使用句柄,HandleScope是Handle的一个栈,javascript类型在c++里面都有对应的类型像String,Integer等,c++通过Handle使用这些类型,通过handle可以使用gc来管理
(3)HandleScope管理Handle的生命周期,HandleScope只能分配在栈上, HandleScope对象声明后, 其后建立的Handle都由HandleScope来管理生命周期,HandleScope对象析构后,其管理的Handle将由GC判断是否回收,对于那种需要return的handle,要用HandleScope::Close转交给上一级HandleScope管理
(4)context_scope,handle_scope都是v8的函数,context_scope意味着进入这个context的范围,后面新建的handle都在这个context下面,直到这个context析构
(5)v8::External把C++的对象包装成Javascript中的变量。External::New接受一个C++对象的指针作为初始化参数,然后返回一个包含这个指针的Handle<External>对象供v8引擎使用
(6)v8::Object这种代表javascript里面的对象

6.CreateEnvironment 首先是诡异的set_as_external,set_binding_cache_object,这些都是宏生成的,暂时理解不了先跳过。然后是uv_check_init(),env->idle_check_handle()这些东东,从注释上看是和profile相关的。另外稍微穿越一下,uv_check_init()初始化的handle是在uv_run->uv__run_check 这里执行的

7.SetupProcessObject 刚读到这里由于对v8的使用方式完全不了解所以还是很晕的,但基本能建立的一个关联是这里在设置process对象,这个process应该和js端经常用到的是一个。比如拿env->SetMethod(process, "_setupNextTick", SetupNextTick)在js里面grep一把立马能看到const tickInfo = process._setupNextTick(_tickCallback, _runMicrotasks)这样的使用,这样就把关联建立起来了

8.LoadEnvironment 这里很让人兴奋,上来就是
  Local script_name = FIXED_ONE_BYTE_STRING(env->isolate(), "node.js");
  Local f_value = ExecuteString(env, MainSource(env), script_name);
这不就是在eval node.js了么?不过还有个小机关,MainSource方法里面引用了一个变量node_native,查看一下它的定义在node_natives.h const unsigned char node_native[] = { 47,47,32,72 ...}
这里js2c.py工具会把src/node.js和lib/*.js转换成字节数组生成node_natives.h
node.js eval出来只是一个函数,所以后面有
  Local arg = env->process_object();
  f->Call(global, 1, &arg);
算是真正启动node.js了

9.v8::platform::PumpMessageLoop(default_platform, isolate) 这个调用一般是执行了自己的script之后,因为v8偶尔会放一些task到前台线程执行所以如果使用default v8::Platform,用户需要自己调用PumpMessageLoop让这些task有机会执行,v8的人建议创建自己的v8::Platform

10.uv_run 启动lib_uv,有一个大循环,在循环里运行uv__run_timers,uv__run_pending,uv__run_idle,uv__run_prepare,uv__io_poll等等,读到这里的时候感觉和java nio的reactor模式比较像但是后面看细节还是有很多不同的,这里就会一直跑大循环直到没有需要监听的时间才退出了。事件循环后面是退出逻辑,做一些退出的回调和资源清理


跟踪一个简单的流程

启动流程有很多没看懂的地方,想继续深入就得有交互,看看nodejs是怎么工作的的了。一开始最好不要弄太复杂的流程,所以我选择跟踪一个写文件的流程

 
var fs = require('fs');
 
fs.writeFile("/home/huying/test.txt","my fs!",function(e){
    console.log('write finished...');
})
//fs.writeFileSync("/home/huying/test.txt","my fs!");
console.log('function finished...');



1.想debug得有一个断点,上面这块代码对应的c++端我根本不知道怎么加断点,所以只能往底层分析js代码。翻了翻fs.js,fs.writeFileSync会调用binding.writeBuffer,然后看到binding的定义binding = process.binding('fs'),看到process眼前一亮,关联建立起来了

2.grep binding我们能找到env->SetMethod(process, "binding", Binding) Binding方法有一句
node_module* mod = get_builtin_module(*module_v)就在这里加个断点,从名字看这一句是在读取builtin的模块,所以就不断f8看这里都bind了哪些module,很快就发现有个node_file.cc的module看起来靠谱。

Binding方法还有一句 mod->nm_context_register_func 内容是node::InitFs,看来module初始化的方法就是它了

3.到这里有两个方向可以探索了,一个是可以继续看node写文件的过程,一个是解答之前的一个疑惑--js和c++是怎样交互的,我选择先了解后者。

InitFs里面有不少env->SetMethod(target, "open", Open)这样的调用,"open"这些名字和js端的调用相同,所以Environment::SetMethod就是c++提供函数给js端调用的地方。接下来又google了一把,找到一篇介绍c++,js交互的好文章,不仅有js调c++的还有c++调java的

http://icyblazek.github.io/blog/2015/02/08/v8-ji-chu-ru-men/

精简摘录一下

(1)设置全局变量给js使用
v8::Handle global = v8::ObjectTemplate::New(globalIsolate);

global->SetAccessor(v8::String::NewFromUtf8(globalIsolate, "globalValue"), (AccessorGetterCallback)globalGetter, (AccessorSetterCallback)globalSetter);

Local source = String::NewFromUtf8(globalIsolate, "var tmpValue = globalValue; globalValue = 21; ");

(2)设置全局函数给js使用
Local globalFunTemplate = v8::FunctionTemplate::New(globalIsolate, (FunctionCallback)globalFun);

global->Set(v8::String::NewFromUtf8(globalIsolate, "globalFun"), globalFunTemplate);

(3)设置一个类给js使用
v8::Local personClass = v8::FunctionTemplate::New(globalIsolate, (FunctionCallback)createPerson);

personClass->SetClassName(v8::String::NewFromUtf8(globalIsolate, "Person"));
v8::Handlep_Prototype = personClass->PrototypeTemplate();

p_Prototype->Set(String::NewFromUtf8(globalIsolate, "sayHello"), FunctionTemplate::New(globalIsolate, Person_SayHello));

v8::Handle personInst = personClass->InstanceTemplate();
personInst->SetInternalFieldCount(1);
global->Set(v8::String::NewFromUtf8(globalIsolate, "Person"), personClass);

Local source = String::NewFromUtf8(globalIsolate, "var p = new Person('Kevin', 'Lu'); p.sayHello();");

(4)c++访问js变量和方法
v8::Local source = String::NewFromUtf8(callJSISolate, "function Person() { this.name = 'Kevin'; } Person.prototype.getName = function () { return this.name; }; var p = new Person();");

v8::Local<Script> compiled_script = v8::Script::Compile(source);
compiled_script->Run();

v8::Handle data_p = context->Global()->Get(String::NewFromUtf8(callJSISolate, "p"));

v8::Handle<Object> object_p = Handle<Object>::Cast(data_p);

v8::Handle getName = Handle::Cast(object_p->Get(String::NewFromUtf8(callJSISolate, "getName")));

Handle value = getName->Call(object_p, 0, NULL);

String::Utf8Value utf8(value);
printf("call js function result: %s\n", *utf8);

到这里就基本搞清楚nodejs c++和js是怎么交互的了,再回去看才c++端启动的时候SetupProcessObject的逻辑就比较容易了

4.process.binding方法在js端经常看到,感觉这是一个很好的debug的点,通过这个点可以比较清晰的了解c++和js两端的交互,所以要重点看一下。

Environment::GetCurrent(args)->static_cast<Environment*>(info.Data().As<v8::External>()->Value()) 这一句上来我就醉了这个info是v8::FunctionCallbackInfo<v8::Value>,这个Data()取的是什么值?我google了半天文档也没找到这个是干啥的,绕了好久才找到关联,在v8-inl.h里面v8::FunctionTemplate::New(isolate(), callback, external, signature) data是第三个参数,而这个external是Environment,所以感觉data是v8给function绑定外部变量的一种方式

之前node启动的时候看不懂的set_as_external(v8::External::New(isolate(), this)) 现在也能理解了,实际上就是把Environment给设置城external变量了。到这里真的感觉看源代码就是在玩一个解谜游戏,不断的找各种线索,建立关联,开启新地图。。。

Binding后面的逻辑还是比较简单的,如果env的binding_cahce里面有moudle就直接返回,否则生成一个新moudle并返回,如果是build的module比如util,会调用到node::util::Initialize 如果是javascript module就直接给module设置name->js source

这里的另外一个收获是,Binding调用栈的前一帧FunctionCallbackArguments::Call(FunctionCallback f),只要是js调c++都会走到这里来,所以这是比Binding更好的一个debug点

5.继续回到之前的那条线看看node怎么操作文件的,先看看同步类型的操作writeFileSync
其实是有三个操作组成的,openSync,writeSync,closeSync
openSync->Open(node_file.cc)->uv_fs_open(fs.c)->uv__fs_work->uv__fs_open->open(系统调用)
writeSync->WriteBuffer(node_file.cc)->uv_fs_write(fs.c)->uv__fs_work->uv__fs_write->pwrite
顺便搜了一把pwrite,这个调用原子操作完成seek和write,多线程操作的时候不用加锁

6.open和openSync的差别从uv_fs_open的POST宏开始,异步流程走到了uv__work_submit

(1)uv_once 调用pthread_once做初始化,threadpool.c->init_once,首先如果配置了线程池大小而且比默认的大就初始化一下线程对应的内存,然后初始化mutex和condition这个和java里面的类似但麻烦不少,然后启动线程池里面的线程。这里注意debug里面显示的线程,在init_once之前有5个线程,1个node主线程,4个v8工作线程,init_once之后又变出了4个node线程

(2)pthread_mutex_lock利用互斥锁加锁,然后把任务放进全局任务队列,调pthread_cond_signal告诉某个线程去读取任务执行任务,这里和java也很类似

(3)这里看到了QUEUE* q参数,从名字上看应该是个队列,应该很简单,但还是花了不少时间才理解。typedef void *QUEUE[2] 定义了void*的数组,而这是在描述一个双向链表结构,QUEUE* q是指向QUEUE的指针,(*(q))[0]是上一个节点,&((*(q))[0])的类型是void**,
(QUEUE **) &((*(q))[0])的意思是把void**转成了QUEUE **,左边再加上一颗*就又把类型转回了QUEUE *所以(*(QUEUE **) &((*(q))[0]))最终返回了QUEUE *就是当前节点的上一个节点,这么曲折是因为要给void*做类型转换。

#define QUEUE_DATA(ptr, type, field)                                          
((type *) ((char *) (ptr) - offsetof(type, field)))

uv_handle_t* handle = QUEUE_DATA(q, uv_handle_t, handle_queue);

#define UV_HANDLE_FIELDS                                                      \
  /* public */                                                                \
  void* data;                                                                 \
  /* read-only */                                                             \
  uv_loop_t* loop;                                                            \
  uv_handle_type type;                                                        \
  /* private */                                                               \
  uv_close_cb close_cb;                                                       \
  void* handle_queue[2];                                                      \
  union {                                                                     \
    int fd;                                                                   \
    void* reserved[4];                                                        \
  } u;                                                                        \
  UV_HANDLE_PRIVATE_FIELDS                                                    \

/* The abstract base class of all handles. */
struct uv_handle_s {
  UV_HANDLE_FIELDS
};

上面这段是要取队列节点上的数据,这里和java双链表最大的差别是java一般是节点包含数据,而这里是双链信息包含在数据里面。QUEUE_DATA的做法是拿ptr(指向void* handle_queue[2])减去handle_queue在uv_handle_t里面的偏移量从而得到uv_handle_t的地址最后做个类型转换。理解了上面两个难点,基本就可以理解QUEUE了

(4)继续看异步工作的流程,把断点加在uv__fs_open(fs.c)上面,就能看到uv__work_submit提交的工作在新生成的线程6里面执行了,worker()->uv__fs_work()->uv__fs_open()

(5)异步工作做完以后会通知主线程,uv__async_send(async.c),这里面实际通知的逻辑是
r = write(fd, buf, len) 居然写文件了,buf的内容就是个1,因为主线程在做io poll,所以马上捕捉到写事件,uv__io_poll(linux-core.c)->uv__async_io(async.c)->uv__async_event->uv__work_done->uv__fs_done(fs.c)  uv__async_io异步回调什么时候被注册的呢?grep了一把,uv__async_start->uv__io_init(&wa->io_watcher, uv__async_io, pipefd[0]),这个模式就很清晰了

这里异步通知的工作方式和java的玩法完全不同,第一感觉是会不会有性能问题?所以继续深入了解了一下,http://www.codexiu.cn/linux/blog/12066/
这里write的fd不是普通的文件,是通过系统调用__NR_eventfd2创建的,是linux提供的一种内建的异步支持实际上是内存中的一个64位无符号型整数,代码读到这个地方隐隐感觉多路复用epoll是个应用很广泛的模式,不仅限于io

(6)再看看open和write操作是怎么连接起来的,
uv__io_poll(linux-core.c)->uv__async_io(async.c)->uv__async_event->uv__work_done->uv__fs_done(fs.c)->After(node_file)->MakeCallback(async-wrap.cc)->node::WriteBuffer->uv_fs_write(fs.c)->uv__work_submit(threadpool.c)->uv__io_poll ...

(7)在AsyncWrap::MakeCallback函数的尾部有这么一段,
  if (tick_info->length() == 0) {
    env()->isolate()->RunMicrotasks();
  }

  if (tick_info->length() == 0) {
    tick_info->set_index(0);
    return ret;
  }

  tick_info->set_in_tick(true);

  env()->tick_callback_function()->Call(process, 0, nullptr);

立马觉得这个和之前关注的技术点tick的实现方式有关联,另外RunMicrotasks这个名字也看起来是一个有故事的函数,所以也了解了一下
http://stackoverflow.com/questions/25915634/difference-between-microtask-and-macrotask-within-an-event-loop-context,不过这里暂时不展开,后面再专注的了解

(8)在我的basicFile.js里面写文件完成以后有一句console.log('write finished...'),console.log经常被用到所以也想顺便看看实现,因为之前找到了很好的js调c++的debug点FunctionCallbackArguments::Call,所以很轻松就找对了地方
FunctionCallbackArguments::Call ->StreamBase::JSMethod -> StreamBase::WriteString


跟踪一个HTTP请求

跟踪http请求是debug网络框架的标准手段,下面是debug用的js代码
 
var http = require('http');

http.createServer(function (req, res) {
    console.log('http function invoked');
    res.writeHead(200, {'Content-Type': 'text/plain'});
    res.end('Hello World\n');
}).listen(3000);
console.log('Server running at http://localhost:3000/');

1.还是要考虑第一个断点加在哪里好,因为前面已经见过了uv__io_poll循环epoll的流程,所以可以先把断点设在这里看看,启动node,然后在浏览器发个http请求进来,果然停在了断点w->cb(loop, w, pe->events)。一路跟踪下去可以得到这样的调用链,uv__io_poll->uv__server_io->TCPWrap::OnConnection,从名字看这是接收了连接。那哪里发起的监听呢?TCPWrap::Listen看起来比较像,加个断点debug一下,果然是的,所以整个流程开始的一段是
js端发起->TCPWrap::Listen->uv_listen->uv_tcp_listen->listen(系统调用)

2.listen调用后面有一句uv__io_start(tcp->loop, &tcp->io_watcher, UV__POLLIN)值得重视,这个函数把&tcp->io_watcher加到loop->watchers里面,loop是lib_uv全局核心数据结构,这个方法其实就是注册一个poll监听的请求,后面的while循环会从请求列表里面把请求一个一个拿出来

3.找到断点怎么加后面就比较简单了,下面是整个流程
TCPWrap::Listen->uv_listen->uv_tcp_listen->listen->uv__io_start()->uv__io_poll->uv__server_io->uv__accept4->TCPWrap::OnConnection->v8->TCPWrap::New->stream.c.uv_accept(初始化client_handle)->MakeCallback(这个应该就是去触发js callback)->v8->StreamBase::ReadStart
->uv_read_start(做stream read相关数据结构的初始化,调用uv__io_start注册read事件监听)->node::Parser::Init()->_http_server.js.connectionListener->node::Parser::Consume()->uv__server_io(next loop)->uv__stream_io->uv__read->read(系统调用)->node::StreamResource::OnRead->node::Parser::Execute->到达状态s_headers_almost_done(遇到\n)->node::Parser::on_headers_complete->node::StreamBase::WriteString->uv_write2->uv__write->write(系统调用)->node::StreamBase::Writev(写入http返回码200)->node::Parser::on_message_complete(into js,其它的on也是在parse过程中调的)->后面就算结束了
这里让我比较纠结的是http parse的过程是在主线程里做的,如果把比较耗cpu的parse放到线程池里面做然后异步通知主线程会不会对cpu使用得更充分一些呢?


了解重点技术点

这个时候对node的代码已经有熟悉感了,所以下一步是集中了解一下比较重要的技术点了。之前提了7个技术点,还有3.nextTick和setTimeout这些是怎么实现的 5.nodejs里面的domain能捕获相关的异常看着很酷,怎么实现的(这个点还是读源代码过程中发现的) 6.epoll为啥那么牛逼 7.了解c/c++的编程风格 这几个问题没有搞定

nextTick相关实现

1.首先可以看看这两篇文章
http://stackoverflow.com/questions/25915634/difference-between-microtask-and-macrotask-within-an-event-loop-context
https://simeneer.blogspot.jp/2016/09/nodejs-eventemitter.html

提交式的异步任务方式分task和microtask,macrotasks有setTimeout, setInterval, setImmediate  microtasks有process.nextTick, Promises

个人理解microtask比较轻,放到当前同步块之后立即执行,堆积起来一起执行,task比较重一般是一个event loop执行一次,两个task之间可以能io,浏览器渲染这些动作,也是把queue执行完,区别是nextTick如果递归会卡住一直执行下去,setTimeout这种如果递归是放到下一个event loop里面跑,所以不会卡死

2.使用https://simeneer.blogspot.jp/2016/09/nodejs-eventemitter.html里面给的代码来debug
 
console.log('<0> schedule with setTimeout in 1-sec');
setTimeout(function () {
    console.log('[0] setTimeout in 1-sec boom!');
}, 1000);

console.log('<1> schedule with setTimeout in 0-sec');
setTimeout(function () {
    console.log('[1] setTimeout in 0-sec boom!');
}, 0);

console.log('<2> schedule with setImmediate');
setImmediate(function () {
    console.log('[2] setImmediate boom!');
});

console.log('<3> A immediately resolved promise');
aPromiseCall().then(function () {
    console.log('[3] promise resolve boom!');
});

console.log('<4> schedule with process.nextTick');
process.nextTick(function () {
    console.log('[4] process.nextTick boom!');
});

function aPromiseCall () {
    return new Promise(function(resolve, reject) {
        return resolve();
    });
}


3.在uv_run while大循环里面执行的
uv_run->uv__run_timers (setTimeout)
uv_run->uv__run_check(setImmediate)

process.nextTick和Promise call(对应Isolate::RunMicrotasks)基本是一起出现,调用的点比较多
比如AsyncWrap::MakeCallback,Environment::KickNextTick等等,基本是处理完一次io就会调用一次

domain的实现

1.因为nodejs有不少当前调用链之外的调用,这些调用的异常无法被当前调用链的try catch捕获,所以用domain统一处理异常,nextTick,timer,event这些是比较典型的使用场景,首先弄一端代码来debug

 
var domain = require("domain");
var d = domain.create();
d.on('error', function(err) {
    console.error('Error caught by domain:', err);
});

d.run(function() {
    process.nextTick(function() {
        fs.readFile('non_existent.js', function(err, str) {
            if(err) throw err;
            else console.log(str);
        });
    });
});


2.这个魔法到底是哪里发生的呢?有两种可能,d.run里面做了某种hack,拦截了异常调用链,或者异常抛到了外面被v8引擎拦截然后用某种事件机制通知回调到domain.run

 
domain2.on('error', function(err){
    console.log("domain2 处理这个错误 ("+err.message+")");
});

try{
    domain2.run(function(){
        console.log('good');
        throw "pig";
    });
    
}catch(err){
    console.log('outter');
}

下面这段代码执行的结果是outter,所以应该是全局捕获了

3.由于domain这个单词有一定特殊性,所以全局grep一下就能找到线索,相关调用链
LoadEnvironment(node.cc)->AddMessageListener(OnMessage)->node::OnMessage->FatalException(node.cc)->process._fatalException(node.js)->domain._errorHandler(domain.js)

process._setupDomainUse(domain.js)->SetupDomainUse(node.cc)->_tickDomainCallback(node.js)


epoll相关

由于能接触到epoll的系统调用了有了一些感性认识,所以想深入的了解一下。但是这里我并不准备debug内核,因为相关的知识和准备还不具备,做这件事情可能相当费时间,所以先在网上找我能理解的文章去建立一些基本认识,等到以后有机会debug内核的时候再去验证这些认识

 1.比较是了解一项技术常用的手段,epoll之前有poll,select,epoll比这两牛在什么地方呢?就从这个角度去搜索,找到了一份不错的文章
http://blog.csdn.net/xiajun07061225/article/details/9250579

2.epoll有三个函数,epoll_create,epoll_ctl,epoll_wait三个函数,epoll_create会返回一个fd
epoll_ctl可以用来增加新的fd到监听fd列表里面
select的api是这样的select(int nfds, fd_set *readfds, fd_set *writefds,fd_set *exceptfds, struct timeval *timeout)
epoll_wait是这样的epoll_wait(int epfd, struct epoll_event *events, int maxevents, int timeout)

所以第一个差别是api导致的,从上面的调用方式就可以看到epoll比select/poll的优越之处:因为后者每次调用时都要传递你所要监控的所有socket给select/poll系统调用,这意味着需要将用户态的socket列表copy到内核态,如果以万计的句柄会导致每次都要copy几十几百KB的内存到内核态,非常低效。而我们调用epoll_wait时就相当于以往调用select/poll,但是这时却不用传递socket句柄给内核,因为内核已经在epoll_ctl中拿到了要监控的句柄列表

3.epoll是事件就绪通知而不是像select那样主动扫描所有监听的fd的集合,所以性能不会随fd增加而下降,这里的核心是epoll维护一个数据结构,其中有一个准备就绪链表,当数据可读时间发生的时候就会把相应的fd放到这个链表里面,epoll只关心缓冲区非满和缓冲区非空事件

4.epoll提供的就绪fd数组使用mmap节省复制开销

5.我猜想的流程是,网卡有数据->内核把数据复制到读缓冲区->触发中断回调程序把socket复制到准备就绪fd列表->epoll_wait返回

6.这里我对读写什么时候会阻塞有一些疑惑,所以顺便找到了一篇不错的文章http://www.cnblogs.com/promise6522/archive/2012/03/03/2377935.html
write成功返回,只是buf中的数据被复制到了kernel中的TCP发送缓冲区。至于数据什么时候被发往网络,什么时候被对方主机接收,什么时候被对方进程读取,系统调用层面不会给予任何保证和通知

编程风格和代码结构

1.基于宏的仿继承,UV_HANDLE_FIELDS<-UV_STREAM_FIELDS<-UV_PIPE_PRIVATE_FIELDS ...

2.基于宏的代码生成 
    set_as_external
    ENVIRONMENT_STRONG_PERSISTENT_PROPERTIES
    先定义V,然后在另一个宏里面去使用V

3.src/*.cc,这里是真正node的c++代码,node.cc是入口和核心启动和装配Process对象以及一些核心的回调方法(给js端用的)。env.cc 类似于一个全局的context用于传输全局参数,比如module_load_list_array=process.moduleLoadList 传送于c++和js之间。xxx_wrap.cc node和lib_uv的桥梁,接js的回调转发到lib_uv的库

4.libuv,uv.h定义uv提供的函数比如uv_run,uv_loop_init等等。win/unix 提供平台相关实现,unix下面有aix.c,kqueue.c,linux-core.c,sunos.c等等

重温启动流程

第一遍看启动流程的时候还有一些地方没太看懂,现在重新复习一下同时也可以把所了解的知识串起来

1.uv__signal_global_once_init 调用uv__make_pipe创建一个管道fd[3,4] (0,1,2是标准错误标准输出这些)属性是uv__signal_lock_pipefd,这对pipeline是用来实现加锁信号量的读pipeline的一端被锁住直到读出合法的值,写的那一端写入*解锁

2.uv_loop_t的属性初始化,loop->timer_heap(timer堆) loop->wq(工作队列) loop->active_reqs(活跃请求队列,active_reqs是个数组联想一下之前看过的QUEUE的定义)  loop->async_handles(异步handle列表,handle类似于http session这种概念,相当于开了一个异步通道,然后这个异步通道上面可以持续有异步请求发生,而request是短暂型对象通常对应handle上的一个io操作,request会用data属性传值,但是多个async handle可以共用一个async fd) loop->nfds(watch的fd的数量) loop->watchers(存放uv__io_t和对应的fd,这个非常非常核心,就是epoll监听的文件描述符列表) loop->pending_queue(发生连接错误或者写结束之类的会把请求扔到这个队列) loop->watcher_queue(待注册epoll的事件会放到这个队列里面) loop->closing_handles(关闭事件的handle列表) loop->signal_pipefd(信号处理的pipeline的一对fd) loop->backend_fd(epoll create出来的fd,后面的epoll调用都依赖这个fd) loop->emfile_fd(EMFILE进程fd用尽,uv__emfile_trick备用一个fd,如果fd用尽就关闭这个备用的fd,然后就可以accept新的连接,accept的同时关闭并告诉客户端fd用尽的状态准备fd的方法是只读的方式打开/根目录)

3.uv__platform_loop_init 调用epoll create初始化fd=5

4.uv_signal_init 创建管道loop->signal_pipefd[6,7] 初始化loop->signal_io_watcher结构,调用uv__io_start把signal_io_watcher加到loop->watchers里面去后面注册epoll就会注册上去。 loop->child_watcher是一个uv_handle_t也会在这里初始化

uv_signal_init 类似这种初始化方式 处理的是linux信号量这种,uv_signal_start 真正启动信号处理,uv__signal_register_handler和uv__signal_handler把流程打通,调用流程是 信号量产生->信号量callback->write pipeline->select on pipeline另一端->调用目标函数

5.初始化读写锁cloexec_lock,打开文件的时候会读这个锁,fork进程的时候会写这个锁。初始化wq_mutex,这个主要事线程池用

6.初始化wq_async,这个也主要是threadpool.c在用,异步io的逻辑主要从这里走。首先通过系统调用__NR_eventfd2创建一个fd,http://www.codexiu.cn/linux/blog/12066/这不是普通的fd,是linux提供的一种内建的异步支持实际上是内存中的一个64位无符号型整数,这里返回的fd=8,调用uv__io_start加到loop里面。后面有一大堆的goto来处理初始化失败去关闭资源,没有exception有点杯具啊


回味一下node的编程模式

























引用一下这张比较经典的event loop的图片,node其实就是在一个主循环里面把大多数逻辑全部做掉。回想之前的http的例子,启动的时候其实就是注册了一大堆listener,然后在poll里面触发事件,由事件触发一长串调用直到写response成功以后返回,然后主线程才能进行下一次poll,而各种nextTick,setTimeout其实就是把任务放到queue里面然后找各个事件函数运行的间歇期去运行。再回想之前写文件的例子,这个例子是真异步,第一次操作在uv__work_submit之后就返回给主线程,主线程这时就继续跑poll,然后异步线程通过uv_async_send去通知poll然后执行之前的异步回调。所以感觉异步就是把程序切成一帧一帧的,帧和帧之间的关联记在某个数据结构里面,然后主线程就是一帧一帧的去跑这些程序帧

node还有一个感觉比较清爽的是把io,async,signal等等用统一的poll循环的方式给表达出来了,编程方式很统一,有点优雅的感觉。

学到了什么?

终于要结束这一段艰难的冒险了,应该还有一些精彩的地方没有浏览到,但那些地方没有和我现有的知识建立关联,而穷举法是效率比较低的方式,所以暂时就到此为止了。总结一下这段冒险自己的收获吧。

1.对node有了一个比较深刻的理解
2.node require元编程的trick令我印象深刻
3.node domain对异步异常的处理方式也很值得学习
4.大致了解了v8引擎的用法
5.对libuv有了一个比较深刻的认识
6.学习了c里面一个比较精致的双向链表实现,同时学会了看复杂的多重指针
7.对信号量,管道,异步这些的理解深了一层,同时学到了一种比较优雅的统一处理这些事件的方式
8.对异步编程的本质理解更加深刻了
9.重温了c和c++的一些语法,对c/c++大型程序如何组织也有了一点认识
10.对epoll的理解深刻了一些
11.学习了一些基本的linux系统调用
12.了解了pthread的基本用法
13.对nodejs积木式创新的模式印象很深刻,同时联想到了weex

这样算下来收货还是挺丰富的,没有白投入这么多的时间

2016年8月16日星期二

分布式和大数据方案备忘

扩容


  1. 分布式缓存(tair) 一致性hash存储,2台扩3台,a,b,a,b,a,b -> a,b,c,b,a,c 只用全量迁移2台机器的数据,成本最低
  2. 分库分表数据库 (1)先做全量迁移并记录迁移之前老库的位点 (2)全量迁移是整表迁移,这也是要做分表的原因之一 (3)全量做完以后开始做增量 (4)增量追上以后把路由规则切到新库
  3. kafka 两种扩容,如果分区开的多只扩机器不扩分区这种可以做到自动迁移,等kafka新机器的isr追上进度即可提供服务。如果扩了分区数,为了保证顺序,必须停写,然后等消费端消费完毕才能扩容
  4. rocket mq扩容,目前对于顺序消息必须停写等消费端消费完毕,因为rocket mq的高可用方案是双写,新机器加进来以后没有isr的方案帮助新机器拿到队列里面的数据
  5. 分布式文件系统 有name server存在,name server会做动态负载均衡,所以新机器加进来即可自动完成扩容和后面的数据均衡
  6. pulsar扩容,在kafka的基础上走一步,把kafka的segments加上元数据管理让一个分区的segments能分散在不同机器上就能做到全自动的迁移和扩容,pulsar的另一个重要特点是存储计算分离原来的broker拆成不存数据的broker和存数据的bookie。但是zk的依赖好像依然是个问题。

热点解决方案

tair hot ring

  1. https://yq.aliyun.com/articles/746727
  2. 基本思路是hash冲突时尽量减少爬链表的过程,把最热的放到最前,这个移动过程可以通过采样来发动
  3. 如果移动链表元素动作较多,并发冲突较多,所以采用ring结构只用移动head指针即可
  4. ring结构转一圈后无法判断是否跑完了,所以要对key排序,这样找到比当前key大的就说明再也找不到了,为了加速比较过程,还在key前面加了tag
  5. 扩容的时候根据范围拆分,原来tag的范围是0-t,拆分成0-t/2和t/2-t,header也分成两个,也就是个标准的rehash手段

hash+range二级散列方案

Hash(a,b)=Range(Hash(a),Hash(b))
有区分度而且保留了一定的范围查询的能力,https://zhuanlan.zhihu.com/p/424174858

一致性

  1. 健康检查类型的数据传播,仅要求最终一致性,所以可以用gossip的方式传播,实现简单
  2. 配置型信息存储,顺序很重要,所以需要用paxos或者raft类型的协议存储


megastore

  1. entity group类似于分表,但是每个数据中心的entity group是相同的。一个entity group内部是强一致事务,跨entity group可以用2pc做强一致事务但推荐的方式是用消息队列做弱一致事务,论文里描述的一个场景是a邮箱给b邮箱发邮件
  2. 高可用性是通过在多个数据中心做paxos写来实现,这种方式的优点是没有master,省去可很多运维上的麻烦。写数据的流程是(1)读取WAL的最后提交位置(2)把要写的数据攒成一个log entry(3)paxos提交(4)把数据合并到big table里面的数据行和索引(5)清理临时数据
  3. 存储和查询基本按关系数据库的做法,但是有主从表的概念,比如user和photo,user是主表photo是从表,两个表相关联的数据在物理上也存在一起
  4. bigtable可以把同一个cell的数据根据时间戳多版本存储,megastore通过这个特性实现mvcc
  5. 有一个coordinator的角色负责记录本地所有paxos写成功的表,通过这个信息尽量保证本地读。这里细节没完全懂,既然是多写,为啥还需要这个角色?
  6. paxos write做了优化,每一轮写都会选一个leader做裁决,选leader的原则是leader所在的区域发起的写最多。具体流程是(1)找leader接受要写的数据如果接受就跳过步骤2(2)选一个最高的sequence,并且把要写的数据更新到最新状态(3)让其它的replica接受更新,如果多数接受就成功,否则重做步骤2(4)invalidate coordinator(5)所有replica把数据更新到本地bigtable

基于json的架构1


1.存储基于mysql,每个表只有两个字段id和json块

2.在mysql里面写,在搜索引擎里面读,mysql和搜索引擎基于binlog同步

3.mysql和搜索引擎上面有一个统一数据层负责读写的分发

4.所有领域对象都是json表示,有一个配置中心负责配置规则做校验和搜索引擎的索引字段

5.服务层上面有一个路由和校验层,根据调用的json串做路由

6.读写延迟是个比较大的问题,无法一致性读导致复杂的业务逻辑可能会比较难做,对数据同步工具要求也非常高


uber schemaless


1.存储基于mysql,cell是基本单元可以是一个json比如乘车基本信息这样的,row_key,column_name,ref_key决定一个cell。一个row可以有多个cell而且不用事先指定。整个表类似于一个三维hashtable row_key->column_name->cell{key1:value1, key2:value2}

2.对cell的操作是基本操作,只能insert不能update或delete,对cell的insert是高版本ref_key覆盖低版本ref_key这样完成update。操作基于cell级别可以有效减少写入数据体积

3.分片存储,每个分片是一个独立的mysql数据库,每个库有一张表存cell实体,每个二级索引用一张表存储,还有一些辅助表。每个cell是实体表的一行,每一行有added_id,row_key,column_name,ref_key,body,created_at,其中added_id是给后面的trigger排序用的,body数据存成压缩的json。这种存储结构是拿多个mysql row拼成一个schemaless的一个逻辑row

4.trigger类似于精卫,有failover,重投这些行为,通过trigger实现异步事件总线,完成类似于行程结束之后付款这样的操作

5.二级索引可以建到cell里面的属性上面,而且支持联合索引,这个感觉是靠索引表实现的。索引也要求有分片字段,索引和cell可能不在一个分片,只做到最终一致,有可能读到旧数据

6.cell存储根据row_key分片,每个存储节点一主两从,尽量读主节点,主从异步复制,读主节点失败会转去读从节点。写主节点失败,会随机找其他分片的主节点写入,但是会告诉客户端这个写入目前不可读,这个叫缓冲写入

7.缓冲写入,写主分片主节点之外会随机挑一个分片的主节点做备份写,在备份分片上的写就是缓冲写,如果主分片主备复制完毕,从节点上出现了新写的数据,worker机器就会把备份分片上的缓冲写数据删掉。相当有趣的解决mysql主备复制问题的方案

8.这个方案有高可用性,扩展性也不错,查询也可以跨维度,感觉比前面的基于json的方案牛逼


缓存使用进化


1.tomcat服务端使用业务对象缓存,web框架对静态页面使用page缓存,浏览器和nginx缓存图片和js这些。缺点是还要走servlet这些,java处理耗cpu,容易到上限

2.动静分离,nginx,缓存服务部署在同一台机器上,静态的数据走缓存,动态的数据由客户端js做csi回源。请求由lvs按一致性hash打进来提高命中率。缺点是管理麻烦,不同的机器间没有互通。

3.统一缓存层,缓存服务单独部署集群,回源采用esi,由缓存服务器找java server回源,这里应该是走比较快的协议。优点是管理方便,而且缓存数据统一互通。失效策略可以有定时,根据tag计算,异步业务消息和监听binlog

4.cdn,先通过dns做路由到最近的cdn节点,cdn如果拿不到数据就找统一缓存层,统一缓存层对动态数据做回源



hbase,lsm tree,列存,kv-store,sstable,skip-list,bloom-filter


1.kv-store是键值对存储,sstable是一种实现方式,特点是按key排序存在内存或者文件里

2.sstable如果没有其它措施读写都会很慢,所以首先有了lsm tree,内存里面L0,然后文件里面L1,L2...一般越来越大. lsm tree一般是写内存,超出大小以后批量merge到磁盘,所以写性能不错

3.说是lsm tree但也可能具体的tree是skip list, skip list加上随机化层数后平衡性不错而且并发的版本加锁比较少,缺点是比较耗内存。skip list一般放内存中,磁盘上还是btree比较好

4.随机查询还是慢,所以用bloom-filter来加速在hfile里面的寻址,通过几个hash函数来快速过滤掉一些不可能有目标数据的hfile

5.hbase的存储建立在上面这些之上,region是分布式存储的最小单元,一个region只会在一台server上,region如果大了就会分裂,region在hdfs上面是个目录。每个列簇对应一个store,一个store里面有一个mem-store和一些文件,这些文件也放在一个目录里面,列簇里面的内容按rowKey排序存在一起,region里面可以有多个store

6.rowKey col1 col2 col3应该是存成 rowKey col1, rowKey col2, rowKey col3 这样

7.hbase称作面向列的存储,存的是列族,但是一般实践都是只有一个列族,列族内的列可以随意修改所以hbase实际上还是行存储。hbase的优势是自动分片,按rowkey查询快,写入快,只能算高级kv数据库。hbase的version一个是确实有场景用得到多版本的数据,另一个是可以用来做mvcc

8.hbase的事务仅仅是针对单行,对于不同行以及单行分布在不同region的情况都无法加锁。hbase有单行的行锁,而且对于单行的读和写可以通过mvcc保证并行度

9.hbase这种存k-v的叫面向列的存储,这种存储的好处是本质上是k-v store可以很灵活的动态添加列,某些列可以很稀疏不占空间,而且可以做lsm-tree结构。kudu这种光存value然后靠偏移量去定位,优点是可以做高压缩比,范围scan很快,这种结构一般要求schema固定有了类型才好算偏移量选择压缩算法,这种结构感觉做lsm-tree也超复杂

hbase row-key

1.hbase region拆分方式是按row-key的range进行拆分
2.row-key太长的缺点是消耗内存和磁盘空间,sstable方式决定
3.row-key简单单调递增会带来当前region读写热点问题以及之前拆分的region很少有机会再被访问的问题
4.热点的解决方案(1)预分区(2)随机散列,缺点是再也无法范围扫描,根据row-key的部分条件扫描也无法弄(3)加hash运算出来的前缀,优点是还可以范围扫,缺点是每次扫都是多机器而且代码稍麻烦
5.row-key可以存字段信息,并且当做查询条件,比如<userId>-<date>-<messageId>这种组合然后查userId 00005 的数据,支持左前缀。这种方式感觉消耗存储而且row-key设计的不好以后不好改
6.还有一种比较诡异的玩法,因为hbase能随意添加列,所以对于用户收到的邮件这种聚合,可以一行,每封邮件一列但是列名不同,通过列名做索引

Cassandra

1.数据存储模型也是bigtable的模型,map<row-key, map<column-key,value>>
2.memtable+sstable+commitlog
3.分片方式是一致性hash,节点的状态是通过gossip协议点对点传播,所以没有了master-server这样的角色,少了单点问题,也不会有hbase region分裂时候的可用性问题,tidb通过raft+状态机解决分裂时的可用性问题
4.高可用的方案是wrn,优点是比较灵活可配置
5.cassandra没有使用vector clock,原因在于本质上cassandra是kv型数据而dynamo是文档型数据,所以cassandra可以用update a=xxx where b=yyy这样的方式单独的更新属性,这样大多数merge就可以避免了

Dynamo

1.大模型上和Cassandra非常相似,也是一致性hash分片,节点的状态也是gossip协议传播,也使用了wrn
2.高可用上面有两招,临时的不可用会找另外一台机器做暂存,等恢复以后再把数据拷贝回来。如果机器挂掉时间比较长,会异步数据恢复,副本之间的数据不一致快速检测和恢复使用的数据结构是Merkle Tree
3.vector clock,某些节点宕机恢复以后或者网络发生partition的时候可能出现版本分叉,比如版本1的基础上出现了版本2和版本3两个不同的修改,客户端读到这个分叉版本的时候需要自己去裁决选择一个版本或者合并,然后再把新版本写回去
4.Dynamo的特点是比较简单轻量,最终一致

td-sql

1.要解决的是mysql异步主备同步可能丢数据的问题
2.最原始的是主挂了停用,捞业务日志补偿备
3.k-v存储加黑名单,主挂了只是没有同步好的黑名单不可用
4.td-sql,一主两备为一组,写主只有binlog到达一备才算写成功,主挂了按binlog最后更新时间选新主

Elasticsearch

1.任意维度的联合查询比mysql强大很多很多
2.term-index(前缀树)->term-dictionary(排序)->position-list
3.多条件合并也就是多个position-list合并,可以通过在position-list上面建skip-list然后合并skip-list,或者直接内存bitset合并
4.skip-list的一种压缩方式是记录delta,bitset的一种压缩方式是模一个65536然后记录商和余数
5.在聚合查询上有和mysql类似甚至更好的能力,因为es又引入了doc_values

6.doc_values是正排索引+列存,比如查年龄18岁的平均身高,先用倒排查到18岁的所有doc_id,再拿doc_id到正排索引里面去匹配然后计算,因为正排是列存的而且还压缩了的所以一次性load进内存进行比较,这样想着感觉好牛,还没完全想明白,要比较详细的了解一下底层的数据结构。通过doc_id拿doc_value的过程是fetch过程,因为要从一堆文件里面定位doc_id所以这个过程其实也很耗时间。

7.es5之前数值索引都是先转换成string然后用前缀fst的方式索引查询,这样的问题是range特别慢,在前缀树上面找range要递归,es6使用了bkd-tree,也就是lsm-tree版的kd-tree,不用b-tree是为了兼容数值型和多维类型的查询。kd-tree的速记方式是二叉树,第一级按第一个维度排序,第二级按第二个维度排序。。。查找的时候也是一个一个维度拿来查

8.es的多条件查询是取交集的形式可能会比较慢,如果配合mysql的查询扫描方式会有很大提升

9.lucene有内存buffer和磁盘segment(可以看成一个独立的子索引)类似lsm tree所以写入高效,但是不能实时查询是因为要给segment建倒排索引,这一步是buffer刷成segment的时候做的所以只能准实时,通过这个lsm tree可以做id到doc的查找

10.es扩展了lucene加了一些系统字段,_uid主键,_version,_seq_no,_source来支持update,_field_names支持稀疏列判断,_routing数据分布,利用的是lucene本身的能力很cool
11.es有个transLog有点像metaq,持久性保证需要每次请求刷盘,如果要实时性可以从transLog查数据,查transLog只能byId,根据倒排索引查询只能准实时。log写放在内存写之后,一个是防丢数据一个是实时查询
12.lucene不支持update,es支持的方式是在内存里面merge成一个新doc,然后add到lucence里面去,其中version起乐观锁的作用
13.海选,粗排,精排,海选一般是根据记录某属性做简单分层,粗排一般做简单算分排序召回1000到2000数据,精排可以有比较复杂的运算做出最终排序



TIDB

1.sharding策略有两种range和hash,range优势是范围扫描和扩容缩容比较好做,tidb使用的这种,缺点是顺序写热点。hash没有顺序写热点问题但无法做范围扫(不过实际业务中的一些加了分库键的范围扫还是可以的),还有就是扩容缩容比较麻烦容易导致较大抖动

2.tidb使用的是多raft-group,就是说每个shard(region)是一个group

mvcc & read skew, write skew

1.time oracle为了解决外部一致性问题,试想如果没有中央协调器,怎么能做到准确的读到一个事务的完整状态?很可能部分节点写完成了,部分节点还没有。而节点很多的分布式环境一个中央协调器又会成瓶颈,所有就需要一个全局的事务id来界定事务所以有了time oracle

2.write screw 两个不同的银行账户各有100块钱,约束是总钱数必须不小于100,两个事务都是先做检查然后取100,mvcc并发状态下两个事务都不违反约束又不会有mvcc冲突。stm内存锁的解决方案是进入事务前都把所有的资源都加个读锁,但是分布式协调器的情况无法统一加锁,所以才有了google的time oracle的方案  

3. percolator没有对read加锁,所以只是一个snapshot isolation,而couchdb这样的对read加锁实现的是serializable snapshot isolation隔离级别,但有性能代价

4.read screw r1[x], w2[x], w2[y], c2, r1[y], c1 or a1

5.引发的问题称作事务的外部一致性

时钟策略

1.逻辑时钟,每台机器维护本地序列号,如果一个事务跨机器从机器1传递到机器2,机器2会比较机器1传过来的序列号和本地的最大序列号取最大的,不同机器的数据只要在事务中形成交叉就相当于建立了逻辑关系。这种方式的缺点是不同事务如果没有数据交叉可能序列号的gap会很大。

2.hlc,物理逻辑混合时钟,在物理时钟不变的情况下新事务只增加逻辑时钟,如果物理时钟变了就跟着最大的物理时钟走。这样相当于用逻辑时钟来校准物理时钟,但不同事务又能通过物理时钟比较大小,但是hlc也无法完全做到不同的没有数据交集之间的事务之间的因果律。

Percolator

1.Percolator是为了补充map-reduce全量离线计算的不足而产生的增量索引系统,首要的优势就是增量计算时效性好。增量计算和全量计算的区别在于全量计算丢弃之前的所有计算结果全部重来,而增量计算会根据新的数据去查找之前相关联的计算结果,然后在之前的计算结果上做局部运算。计算的可加和性是增量计算是否便于实施的关键

2.为啥需要事务?爬虫从多个url抓到了相同的内容,只需要将pagerank最高的url加到索引,但是每个外部链接也会被反向处理,让其锚文本附加到链接指向的页面,指向复制品--pagerank低的那些页面的链接必要时也会被修改指向pagerank最高的url。这里就涉及到正向索引,反向链接锚文本和指向复制页面链接等几个状态的修改,所以需要一致性事务。还有一个可能的原因是一个文档的pagerank,内容hash这些可能不是放在一张表里面而是分开存放的

3.percolator的问题域特点是海量数据和延迟容忍度相对高,延迟容忍高所以对脏锁清理可以比较懒惰

4.percolator也是2pc事务但是和传统的2pc事务最大的差异点在于没有中央协调者,而这个特点带来了极大的可伸缩性

5.能做到无协调者的几个关键因素是,mvcc快照隔离级别,bigtable的mvcc单行事务和可扩展性存储,轻量的统一时间戳提供者

6.percolator的表的一行内容是,C:data(实际存储的数据),C:wirte(控制数据可见性是否commit),C:lock(事务用的锁),C:notify(通知机制用,避免全数据扫描),C:ack(防止重复通知)

7.mvcc和乐观锁的区别在于mvcc可能是多行的一个统一的状态,比如事务1写了行a然后去读取行b,这中间事务2写了行b,那事务1应该是读取事务2写之前的版本,所以这里就需要有时间戳比较或者版本号。而乐观锁只是锁定一行就行了。单机版的mvcc时间戳比较容易做,而分布式的mvcc就需要time oracle

8.论文中代码需要注意的几个点(1)被实现的Transaction是分布式事务,调用的bigtable::StartRowTransaction是bigtable单行事务 (2)oracle.GetTimestamp()是取统一时间戳
(3)注意区分代码中读写数据的时候"lock","data","write"几种不同的标签 (4)Get方法最有意思的是先去读一下当前row的锁,如果读到了就可能清理过期的锁这也是因为没有协调者所以只能靠读的时候清理锁 (5)数据提交过程类似2pc分两步,第一步Prewrite写入数据和锁信息,这一步可能会发生冲突而失败,第二步真正让数据可见也就是写入"write"列并且会把锁擦除。(6)注意每次写入的时候会选一个Primary Record,这个玩法类似传统2pc Primary写入成功作为如果发生异常后面的客户端判断这个事务该回滚还是该提交的标志。(7)提交的时候擦除锁是一个一个进行的,所以Get()方法碰到锁的时候有一个等待,这样就不会读到部分成功的状态

9.time oracle通过定期分配一个时间范围放在内存的方式避免磁盘io,worker也有定期批量从time oracle获取时间戳的机制,这样time oracle单机tps近两百万,这个地方无法得知更多细节但是思想大致能懂

10.数据触发器这点上使用的是弱一致的通知语义避免锁占用时间太长,而实现通知语义使用的是线程范围扫描,有两个提高性能的点,一个是使用notify列减少扫描数量(因为有变更的毕竟是稀疏的),第二个是后面的线程扫描和前面的线程冲突的时候会再做一次随机选择避免公交车效应,这里对细节不是很清楚

11.Percolator处理一个文档需要50个bigtable操作,rpc多,所以使用了很多微batch等待打包rpc请求,还有通过对一行数据的预读的方式提高性能

spanner

1.和percolator最核心的区别是spanner是全球分布的而且是一个数据库,所以不能用单一的time oracle,数据时效性也比较重要所以必须有事务协调者

2.多个time server就意味着时间差异,spanner使用了原子钟+gps保证time api高可用而且时间误差非常小,多台time server之前会互相做时间校对,而且每个时间客户端也会从多台time server拉取时间获取综合值

3.有了时间误差就有了置信区间的概念,从time api取得的时间是有误差区间的,所以分配时间戳的时候会把误差区间考虑进去从而保证事务之间的绝对时间顺序关系,甚至会采取等待的方式

4.数据高可用是采取带leader的paxos变种实现的,最基本的单位是tablet(类似于数据库表),每个tablet的几个replica形成一个paxos group,replica的数量和分布都灵活可配置。读写发生在多个paxos group就形成了分布式事务

5.spanner牛逼的一个地方是数据分片会很灵活,Spanner底层存储的是有序的KeyValue集合,在数据模型上,细化了directory这个概念。一个directory是连续的一组KeyValue,比如user_id=1的用户和他所有的相片信息就组成了一个directory,多个directory组成一个tablet,但又不像bigtable那样每个tablet里面的主键都是连续的按范围划分的。而是通过placement server做动态调整,把经常一起使用的向物理位置近的地方放,做到尽量优化

6.分布式事务采用的是2pc+mvcc的方式实现,每个事务都会有协调者,协调者的选举应该也是和高可用的leader选举类似。剩下的核心点就是怎么控制time api的误差来做到先提交的事务先可见。之前有个疑问是为啥不用协调者生成事务id,这是因为协调者可能会经常变,而且不同事务的协调者可能不是一个

7.spanner的数据模型还是google最常用的基于嵌套结构的半关系型模型,spanner tablet的存储是类b tree,这个细节会有什么影响还需要后续了解

google dremel

1.第一个核心问题是通过列的方式存储嵌套的数据结构,所以有了r,d的数据描述和有限状态机的数据读取形式
2.r代表当前记录和上一条记录在哪个嵌套深度有公共节点,d代表当前记录walk到根节点需要走过的记录数,d对有内容的节点其实意义不大因为都应该是一样的但是对null意义比较大,能还原出null所在的层级,通过r和d就能完整的还原出整个层次结构
3.只存储叶子节点的数据,读取过程是通过有限状态机在不同列之间跳转读取来还原嵌套结构
4.查询的时候好像也没有什么大招,就是先还原出嵌套结构,然后过滤,但感觉这个顺序可以穿插起来如果是单条件查询感觉还是没有什么好招去加速
5.嵌套的数据结构和关系型一对多多层关联相比,保留了数据关系,跨层关联的时候比如查询Country是xxx的name有多少个,关系型数据库就是三层join,而嵌套型结构需要遍历的数据会少很多,在dremel里面可能是先查一下Country列,然后根据r和d判断一下对应的name就可以了,都不一定要走中间language层。嵌套型的缺点是可能会有大量数据冗余,视具体情况而定

Durid

1.实时olap,亚秒级查询,实时数据导入立即能查到,不支持join,列存,有时间序列数据库的特点,带上时间查询会方便寻址  

2.olap基本理论,列按照职责划分timestamp column(时间戳),dimension columns(过滤或者聚合的维度),Metrics(聚合和计算的基本数值)

3.durid不存基本数据而是在数据导入的时候预先对数据做各个维度的first level roll up即预先做好基本聚合,是典型的molap,这个和工作中遇到的数据应用场景完全吻合

4.durid分实时节点和历史节点,实时节点定期merge到历史节点,查询也可能merge实时节点和历史节点的数据

5.segment是存储的基本单位代表一个时间窗口内的数据,Segment文件名称的格式:dataSource_interval_version_partitionNumber

6.完全的列存,只保存列值不保存row-key,会采用字段表的方式压缩,还有位图信息做过滤筛选

Kudu

1.融合oltp和olap,hbase,cassandra这些大范围scan性能有限,纯列存数据库随机读写性能差。kudu大范围scan性能高,随机读写性能不太差

2.真正的列存而不是sstable,schema固定,数据有类型。有一个隐性的行号列,通过这个行号列来定位列里面的记录

3.分区方式是hash+range结合,兼顾范围扫描和避免局部热点

4.RowSet是最小存储单位,又分为MemRowSet和DiskRowSet,内存中用的是行存储刷到硬盘中使用的是列存储。Kudu的存储是base+delta的形式不同于lsm,lsm永远先写MemStore,Kudu是insert先进mem但是update可能会直接写到disk,lsm的key对应的值可能在多个store但是kudu的只会在一个RowSet的范围内,虽然可能同时在MemRowSet和DiskRowSet甚至delta file里面  

5.一个RowSet只有一个MemRowSet,数据insert的时候都会进MemRowSet,实现是一个并发优化过的b+ tree。RowSet有一个DiskRowSet(base file)和多个DeltaFile(undo,redo record)

6.DiskRowSet是列存,会被切成多个page而这些page通过b tree管理,也用了bloom filter来加速对pk的查找。DeltaStore也分Mem和Disk,会组织一个map,key是(row offset+timestamp),value是rowchangelist。update的时候会先通过bloom filter和b tree从DiskRowSet里面找到记录的offset,然后在delta file对应的offset的位置写入值

7.kudu也支持mvcc以及external consistency,也有类似于time oracle的hybrid time

Calvin

1.解决事务并发冲突的办法是让事务根本无法并发,把所有事务做一个全局排序变成确定性的,有一个关键点是batch,不然无法做依赖排序

2.许多sequencer拦截所有的事务请求然后保存到局部一点的存储并同步到其它sequencer对应的存储,这些存储又被一个中央的meta存储管理

3.很多scheduler并行去执行这些事务,如果所有事务单线程执行当然不会出错但效率太低,我的理解是scheduler会分析这些事务的数据依赖,然后尽量并行它们。这个方案给我的感觉有点类似oz的思想

图数据库  

1.索引存储的方式是用顶点+邻接表的方式存储,所以邻接表拿到顶点列表,以及通过顶点拿边都是O(1)的时间,所以就很快了。而mysql join找相邻顶点列表,以及通过顶点找边都要扫表    

2.比如有一个user表,一个item表,要存储好友关系和购买关系,就会有一张friendship表和一张purchase表。这个存储方式和关系数据库里面自己建表的方式区别不大,但是在查询的时候我理解图数据库是用图搜索的方式推进的而关系数据库是通过join实现的,所以在查询多层关系(>3)的时候图数据库性能有指数级别优势


向量化执行引擎

相关的两个关键字是列存和SIMD,SIMD是多个操作数组成向量一个指令完成多个操作数运算。向量化执行引擎充分利用SIMD,而只有列存的情况一列数据才会很容易load成向量利用SIMD进行操作


KAD(Kademlia)算法

1.这是一个分布式hash的算法,可以参见文章https://www.jianshu.com/p/f2c31e632f1d  

2.首先是存储内容按hash值尽量均匀分配存储到多个节点,并且有冗余备份

3.算法的核心是怎么寻址,地址不可能在每个节点全量存储,所以每个节点只能存部分地址,这种情况寻址肯定是跳多个节点。整个算法的亮点是发明了一种xor距离的概念,xor距离带来了距离分层2^0, 2^1, 2^2 ...在相同层次k的节点之间的xor距离必定小于2^(k-1)。所以这样带来的结果是最多k次寻址就可以找遍2^k个节点,而且每个节点的地址表理论上存k个地址就够了


IPFS(星际文件系统)

1.对标http,http是中心化存储,ipfs是真正意义上的分布式文件存储。而且由于merkledag的数据结构,可以做到分布式环境下的去重。  

2.个人理解的核心是把文件拆成块使用merkledag数据结构存储,可以起到文件寻址,去重,防篡改等类似区块链的效果,然后对每个块使用KAD算法做分布寻址。  

3.在应用上增加激励机制来鼓励分享这样会有很多应用的想象空间。


Succinct Trie

1.高压缩trie tree,succint数据结构是一种高压缩数据表示方法  https://www.jianshu.com/p/47a61caa6490

2.succint主要是用位图的方式表达了树的结构关系,核心在于rank1(0)和select1(0)两种操作分别代表第position[0 ... x]中1(0)的个数和第x个1的position。对树的编码方式是节点r有0个子节点就是0,有1个子节点是10,2个子节点110 ... 然后按广度遍历的方式放置成一个bitmap
参考文章https://www.jianshu.com/p/36781efac8e9  

3.树的子节点或者父亲节点都能用rank和select两种操作简单表达出来

4.Fast Succinct Tries分sparse和dense两种,dense一般在trie上层因为上层一般比较满,sparse在下层。dense每个节点,1层放所有字符比如a-z的bitmap,2层bitmap放有没有子节点,3层bitmap放是否是前缀结束,4层是合法前缀对应的值。查询的时候2层用来做tree导航。sparse结构每个节点,1层放字母的二进制,2层用一位表达是否有子节点,3层表达是否是父节点的第一个子节点,4层放合法前缀对应的值。2,3层用来做tree导航。

5.Fast Succinct Tries点查比bloom filter差,范围查好很多


Fractal tree(toku db)

1.mysql(innodb) pk 顺序插入快,随机插入慢。原因是随机插入可能造成节点分裂写放大,存储不连续碎片化。
2.mysql有插入缓冲的解决方式只能针对非主键,大致做法就是增加一个缓冲区,批量刷出。
3.fractal tree的解决方案,结构上类似b+ tree,但每个节点都带msg buffer。插入操作的时候先找到某个节点对应的msg buffer然后放进去就ok了,后台线程异步刷出去。怎么找这个节点的细节还不太清楚。查询的时候需要合并整个查询路径上的msg buffer上面的信息,所以会比较慢。

RAFT

1.通过广播消息+timeout选master,拿到多数票胜。
2.只写master,写了之后记wal日志然后广播到其他节点,拿到多数日志算提交记录commit日志并返回,否则回滚。这里wal日志和commit日志也可以合并。
3.新节点通过master的commit日志来追进度达到一致状态,也可以通过定时的snapshot快速追状态。

存算分离

1.解决成本问题,比如双11只是计算突然飙升但是存储没有太大增加,扩容性价比太低
2.存算一体数据库扩容故障恢复要搬数据(多数情况),恢复时间长
3.存算分离可以使用比一体式方案大很多的存储,大幅度减小分布式事务分布式查询的概率
4.mysql存算分离不是简单的把本地存储换成云存储,需要有存储引擎的较大改变,比如Aurora,polardb,cynosdb。但似乎存算分离的数据库还是单实例数据库,并不能算分布式数据库。
polardb的分布式应该是在存算一体的基础上附加了分布路由和事务
5.这样做的原因是(1)大量存储空间浪费,page+各种log*主备实例数*云盘三副本(2)本地io变网络io各种性能消耗,带宽消耗
6.db即redo log,只保留一份redo log,其他所有的数据都可以从redo log恢复。写入操作只要redo log落存储层就算成功。redo log被分成多个segment来实现扩展性,log和data遵守一致的分片策略。存储层有能力从log恢复出数据页(传统数据库是计算层做的),计算层读取缺页的时候可以直接从存储层读取
7.计算层依然分主从,但主从都从统一的存储读取。计算层完全无状态,可以快速恢复。
8.CynosStore存储的特点,第一是非对称读写,(写的是日志、读的是数据)。第二是异步写,同步读。第三并发写入、日志串行化。第四是支持两层,块设备层、文件系统层。第五能接入任何基于日志先写的存储系统。第六是存储分布式化

RoaringBitMap

1.主要为了解决bit map稀疏时浪费空间较大的问题,比如32位数的bitmap存储是512M,但如果数很少会比较浪费
2.Roaring Bitmaps把32位数的空间划分成2^16个桶,每个桶最多可以存放2^16个数,刚好覆盖整个空间,桶里面存的是数字的低16位
3.保存一个数的时候用数的高16位作为桶的编号,如果桶内的数较少少于4096,就存成array,这样比较省空间,比如只有4个数就只用存4*2个字节,array动态扩容。如果数较多就会存成bitmap,2^16个bit是8KB的空间
4.如果存满是65536*8K=512M,和bit map相同。参考不深入而浅出的Roaring Bitmaps原理
5.64位不能简单像32位那样搞2^32个桶因为一个2^32的数组占用空间已经很大了。一种做法是前32位用红黑树存,更好一些的做法是前48位用art tree存,后16位还是放在一个个桶里,art tree可以认为是可变长的tried tree,每层的节点数可以是4个,16个,48个或者256个。

Z-order Index

  1. 参考Z-order是如何提升查询性能100倍的
  2. 解决的问题还是b tree多维度索引必须前缀匹配的问题,z-order的思路是把多维的信息压缩成一维的顺序,而且每个维度平等,尽量保证局部性
  3. 压缩下来的效果的例子(a1,b1),(a1,b2),(a2,b1),(a2,b2),(a2,b3),(a2,b4),(a3,b1),(a3,b2)...
  4. 二维信息压缩的方式,将x和y表示为二进制以后,每个bit位相互插值即可产生Z-Address
  5. 这样做的好处是多个维度查询都能获得不错的性能,但缺点是单一维度没法获得极致的性能了,这可能也是clickhouse没有使用z-order的原因
  6. 本质是一位有序路径上尽量保证高维度上的局部性,Z-order的Z字形路径其实还不完美,还有希尔伯特曲线路径等等

向量检索

  1. 向量索引对LLM有用,因为LLM context不能太长(内存,计算复杂度,长期依赖),向量化可以储存context知识,相当于LLM的一份外存
  2. 一种比较优秀的向量检索算法是scann
  3. 算法的主要思想是,近似计算,信息压缩,分片查询
  4. 首先是Vector quantization,快速理解可以看矢量量化,把高维(d)向量拆成M块,对每一块做KNN选256个中心点(8位),然后把每一块数据替换成0-255中间的一个,这样原来的数据变成了M个byte。原来的内积相似搜索O(Nd)变成现在的O(256*d+MN),M<<d
  5. 同样通过K-means把把空间划分成多个部分,看查询点q处于哪个部分,然后只比较那个部分所有的点就好了
  6. 在这些的基础上scann提出了各向异性量化损失函数提升计算精度,细节没太懂基本思想是和传统的内积损失函数比,增加了和向量方向相同的权重,降低正交方向的权重

XOR Filter&Ribbon Filter

  1. xor filter 比Bloom Filter节省25%空间!Ribbon Filter在Lindorm中的应用
  2. bloom filter的主要问题是空间占用较大1.52n(n=log...)所以有了省空间付出额外cpu的xor filter(1.23n)
  3. xor filter的主储存是m=1.23n个元素的数组,每个元素r位。查找的时候对于key,算三个hash到数组三个位置拿出三个r位的元素做xor,xor的结果用来和footprint(key)做比较,不相等就确认false
  4. footprint函数怎么来的还不太清楚,已有footprint函数和3n个hash函数的情况下构建xor filter的过程大致是每个元素找不和其他元素hash后发生位置冲突的位置,重复这个过程直到所有元素放进数组(如果找不到了就换hash函数重来)。最后插入的item位置直接设置对应footprint函数的值,然后开始反推其他item位置对应的hash值(用已有hash xor footprint)
  5. ribbon filter是在xor基础上进一步优化空间,基本思路是类似的,主要的区别在于三个hash函数的结果是一个向量,filter数组是一个0,1矩阵,filter的过程是计算矩阵乘法。filter构建的过程是高斯消元法(解多元一次方程组)
  6. ribbon filter的filter数组是列存,过滤时不需要一次全部加载可以一列一列计算有短路逻辑
  7. ribbon的r是7而xor是8所以ribbon的存储是1.101,之所以能做到是因为矩阵信息更密集?



2016年8月6日星期六

对区块链的一点理解

读了一篇很不错的讲区块链原理的文章,整理了一下思路 http://o.btc123.com/data/docs/easy_understood_bitcoin_mechanism.pdf



在知道谜底的情况下重新推导一下

1.需要去中心化,所以选择分布式存储,数据量可接受,所以存全量
2.交易数据需要正确性,顺序性,所以存储全部历史,链式存储
3.交易安全保证,使用公匙,私匙协议,这个比较常见
4.需要解决拜占庭将军问题,首先是一致性问题,这个类似于抢占式协议
5.数据伪造问题,这个是很多分布式协议里面碰不到的,所以有了牛逼的工作量证明协议
6.工作量证明协议最重要的是,验证容易,计算困难,计算过程可协作
7.本质上还是一个保证一致性的分布式存储系统

少了一个Merkle Tree

1.这个结构用来在整个链中间验证交易是否存在
2.从最长的链拿到根hash值和该交易认证路径,然后根据路径重新计算根hash值和实际的比较就行了
3.在数据库类型的系统中会有这个快速定位修改

之前对pk的过程理解得不是很清楚

1.主要是为了防止双花,就是说a客户端先给b付了一笔钱,然后又把这笔钱付给c(前提是这个客户端被改过了,不然本地校验过不去),然后创建一个强大的支线链路争取超过第一笔交易
2.但是第一笔钱已经跑了6个区块了,第二笔钱的链无论怎样也追不上,所以不会被系统承认,如果没有跑6个区块那第一笔交易其实没有达成如果被第二笔超过那就只成交第二笔就好了
3.之前有个地方理解不对,分支长链并不会把之前短链的交易都干掉,只会让短链上的工作都转到长链上面来


2016年4月28日星期四

学习一下kafka

学习的代码版本是0.8.2.1

原理

组件图


存储

目录结构

tmp$ tree kafka-logs/
kafka-logs/
|-- huying_test-0
|   |-- 00000000000000000000.index
|   `-- 00000000000000000000.log
|-- huying_test_ultimate-0
|   |-- 00000000000000000000.index
|   `-- 00000000000000000000.log
|-- huying_test_ultimate-1
|   |-- 00000000000000000000.index
|   `-- 00000000000000000000.log
|-- huying_test_ultimate-2
|   |-- 00000000000000000000.index
|   `-- 00000000000000000000.log
|-- huying_test_ultimate-3
|   |-- 00000000000000000000.index
|   `-- 00000000000000000000.log
|-- recovery-point-offset-checkpoint
`-- replication-offset-checkpoint

一个目录代表一个(topic, partition)组合,在代码中是一个Log对象,每个Log包含多个LogSegment,一个LogSegment包含一个log文件(FileMessageSet)和一个index(OffsetIndex)

Log

+----------------+
|offset 8(bytes) |
+----------------+
|messageSize 4   |
+----------------+
|CRC 4           |
+----------------+
|Magic 1         |
+----------------+
|attributes 1    |
+----------------+
|key 4           |
+----------------+
|size 4          |
+----------------+
|content         |
+----------------+
+----------------+
|offset 8(bytes) |
+----------------+
|messageSize 4   |
+----------------+
|CRC 4           |
+----------------+
|Magic 1         |
+----------------+
|attributes 1    |
+----------------+
|key 4           |
+----------------+
|size 4          |
+----------------+
|content         |
+----------------+

OffsetIndex

+------------------+
|relativeOffset 4  |
+------------------+
|positionInByte 4  |
+------------------+

LogManager startup

1.遍历所有的日志目录,如果有.kafka_cleanshutdown就说明是干净的shutdown,跳过当前log的recovery
2.读取recovery-point-offset-checkpoint文件,返回(topic, partition)->offset 映射。其中文件的第一行是版本号,第二行是记录条数,后面是映射内容
3.计算出log目录和对应的recovery point,并用来初始化Log对象。
4.清理工作,清理掉.kafka_cleanshutdown
5.启动三个定时任务,kafka-log-retention, kafka-log-flusher, kafka-recovery-point-checkpoint
6.retention是删除时间过久,体积过大的log文件。flusher是根据时间间隔刷盘,调用FileChannel.force。recovery-point-checkpoint是把内存里的(topic, partition)->offset checkpoint结构刷盘。
7.启动LogCleaner

new Log()
1.删除掉.deleted和.cleaned后缀的文件
2.碰到.swap说明在swap中途server挂掉,这里的swap过程仅仅是一个rename。这种情况对于index的swap直接删除后面从log文件重建,如果是log的swap就先删index再通过rename的方式swap回log
3.遇到.index索引文件,如果没有对应的.log文件直接删除
4.通过.log文件生成LogSegment对象,如果没有索引文件就运行LogSegment.recover重建索引
5.重建索引的过程是以message为单位遍历LogSegment,遍历过程中如果累加的大小超过了index interval就在索引里面记录一下当前消息偏移量(会转化成相对偏移)和.log文件里面的字节为单位的位置。这个过程中会顺便把log文件和index文件末尾的可能因为crash产生的多余的字节给清理掉。
6.开始根据recovery point做恢复,这里碰到cleanShutdown文件就会跳过恢复
7.遍历从recovery point到末尾的LogSegment,调用LogSegment.recover,如果过程中发现多余的字符说明LogSegment非法,就把非法的一直到末尾的全部删掉。
8.最后对所有的index做一遍sanityCheck

Log.append

1.参数是MessageSet一次append一批消息
2.分析校验,拿到是否递增,压缩codec,以及校验大小,crc,干掉多余的字节
3.后面的所有过程都加上Log级别的锁,如果需要设置偏移量,就遍历每条消息,设置上偏移量
4.再做一次消息大小检查,因为前面可能重新压缩过消息
5.如果当前LogSegment大小加上要append的大小超过上限就滚动到新的segment
6.用当前Log的logEndOffset生成新的log file和index file,之前的末尾index文件做一下裁剪
7.用新的log file和index file生成新的LogSegment并加入Log里面的列表,并且把recoveryPoint到之前的logEndOffset的segment加入刷盘队列,刷完了以后recoveryPoint增长到logEndOffset
8.如果增长的字节量达到了需要做索引的长度就在index里面append一个entry,然后FileMessageSet.append
9.如果append以后没刷盘的字节过多就刷一下盘
10.一个有趣的点是这里并没有有mapped file但是在索引的地方用了,而rocketmq是主log文件也用了,个人理解是只有读也比较多的时候用这个才有价值,只是一次append然后一次读意义就不是很大

Log.read

1.参数有一个startOffset,首先定位到offset仅仅比这个低的entry,如果定位不到就直接异常
2.调用LogSegment.read,首先要找到log文件里面对应startOffset的第一个合法的position
3.查找过程先通过index查找offset对应的position,查找的过程是标准二分查找,因为用了mmap所以都是在内存里面找。index因为需要经常的随机读和append操作,所以做了mmap
4.调用FileMessageSet.searchFor,在log文件里找大于offset的合法message的position
5.用同样的方法计算出end offset对应position,在算出需要读取的字节数
6.通过这些信息,返回一个FileMessageSet,这里还没有实际读取,后面会提取出byte message返回

LogCleaner

1.如果对一个消息存在key相同但是offset更高的消息,那么就清理掉这个offset低的消息
2.先找最脏的Log,找的方法是先读取cleaner-offset-checkpoint里面存了每个log的最后清理offset,然后每个log拿最后清理的offset到active offset所有segments的size和log size算一个比率,然后取比率最大的
3.开始清理Log,先遍历所有脏的LogSegment,构建一个内存中的OffsetMap(index),key是消息key,value是消息offset,这个map后面用来判断消息是否是重复的
4.把所有的LogSegment分个组,每个组的字节大小总量尽量接近于某个指定值。尽量把一个组压缩成一个新的Segment,这样做让新的Segment体积尽量平均
5.拿出每个组的第一个segment以它的名字创建log和index对应的.cleaned文件
6.对于组里的每个LogSegment,过滤掉重复的和空的message,其余的全部写入那个新的目标LogSegment
7.新的LogSegment刷盘,然后把旧的segment list全部删除,最后把.cleaned名字改成.log文件
8.为啥会有这种清理策略呢?其实是为了用topic保存位点做准备工作,保存位点的场景key比较固定只需要最后状态,所以用这招很好。如果是交易订单的场景就不合适了

网络协议

request and response

1.首先解析2字节的requestId
2.根据requestId选择kafka.api包里面对应的request类来处理,序列化反序列化协议自包含在对应的request类里面
3.有一个correlationId做request和response的关联
4.拿ProducerRequest做例子
+----------------------------+
|     requestId 2            |
+----------------------------+
|     versionId 2            |
+----------------------------+
|     correlationId 4        |
+----------------------------+
|     clientId               |
+----------------------------+
|     requiredAcks 2         |
+----------------------------+
|     ackTimeoutMs 4         |
+----------------------------+
|     topicCount 4           |
+----------------------------+
| +-------------------------+|
| |   topic                 ||
| +-------------------------+|
| |   partitionCount 4      ||
| +-------------------------+|
| | +----------------------+||
| | | partition 4          |||
| | +----------------------+||
| | | messageSetSize 4     |||
| | +----------------------+||
| | | messageSetContent    |||
| | |                      |||
| | +----------------------+||
| +-------------------------+|
+----------------------------+

Message

+--------------+
|CRC 4         |
+--------------+
|Magic 1       |
+--------------+
|attributes 1  |
+--------------+
|key 4         |
+--------------+
|size 4        |
+--------------+
|content       |
+--------------+

消息发送

1.发送分同步发送和异步发送,参见Producer.send,最终都会调用DefaultEventHandler.handle
2.发送消息之前如果topic metadata的缓存更新周期到了先随机找一台server拉一下数据更新brokerPartitionInfo
3.如果这条消息没key就和以前发送的partition相同或者随机选一个分区,如果有key就通过partitioner计算分区号,发送的broker选择leader broker
4.发送失败的消息会进行重试,有条件remainingRetries > 0 && outstandingProduceRequests.size > 0进行判断
5.ProducerSendThread里面有对异步消息批量发送的机制,通过数量多和超时两个条件判断是否该发送消息。可以选择是否要ack,有个配置项request.required.acks
6.server端KafkaApis.handleProducerOrOffsetCommitRequest,首先如果是offset commit request就转换成对offset topic发送消息的producer request
7.调用Partition.appendMessagesToLeader,把message保存到本地的文件系统。调用replicaManager.unblockDelayedFetchRequests把block住的fetch request处理一下符合条件的给response。delay fetch需要满足的条件在DelayedFetch.isSatisfied,主要是积攒一定量的消息
8.根据ack方式的设置,或者直接给response,或者生成DelayedProduce放到ProducerRequestPurgatory里面等后面满足条件的时候再ack。条件在DelayedProduce里面,主要是检查所有replica的end offset
9.保序和线程模型的问题,服务端是多个线程共同处理请求的,所以如果想保序的多条消息同时被服务端处理肯定乱序,但因为有ack机制,发送端发一条ack之后再下一条,这样在同一时间点就只能存在一条需要保序的消息就不会乱序了。

接收消息

1.ZookeeperConsumerConnector主要处理consumer和zk的交互,/consumers/[group_id]/ids[consumer_id] -> topic1,...topicN 这个节点是consumer注册的临时节点并且存放了这个consumer订阅的所有topic,consumer会监听自己group里面其它consumer的变更,感觉是个坑,当一个group里面consumer很多的时候zk监听会非常的多
2.consumer会监听/brokers/[0...N],感觉也是坑。/consumers/[group_id]/owner/[topic]/[broker_id-partition_id] --> consumer_node_id 记录partition的owner,rebalance的时候会重建
3./consumers/[group_id]/offsets/[topic]/[broker_id-partition_id] --> offset_counter_value 记录消费位点,还是感觉用专用的topic记录比较好
4.consumer启动api,Consumer.create(config)实际上启动了ZookeeperConsumerConnector,启动过程中如果配置了autoCommitEnable会启动一个scheduler定期扫描consumer对应的topic,partition并且提交位点
5. 调用ZookeeperConsumerConnector.createMessageStreams,首先在zk consume group的path下面注册client的临时节点,然后初始化3个zk listener ZKRebalancerListener,ZKSessionExpireListener,ZKTopicPartitionChangeListener
6. ZKSessionExpireListener不监听具体节点,在重连的时候重新注册client临时节点,触发loadBalancerListener.syncedRebalance。ZKTopicPartitionChangeListener监听/brokers/topics/huying_test,发生变化的时候触发rebalance,loadBalancerListener.rebalanceEventTriggered
7.ZKRebalancerListener监听/consumers/console-consumer-84505/ids节点,ZKRebalancerListener启动的时候会启动一个定时任务检查是否发生过rebalance事件,如果发生过就调用syncedRebalance。
8.rebalance的过程先检查broker,如果没有broker就注册一个listener。为了防止rebalance失败导致重复拉取,停止当前consumer上的ConsumerFetcher,清掉线程消费queue和KafkaStream里面的消息。接下来删除zk上的partition ownership信息,和内存里的注册
9.接下来调用PartitionAssignor.assign给当前consumer重新分配partition,有两种分配方式RoundRobinAssignor和RangeAssignor,RoundRobinAssignor先拿到所有的partition和consumer thread,然后分别排序再循环consumer thread去拿partition,但是只返回本consumer对应thread id对应的partition。RangeAssignor也比较简单,比如有10个partition,有5个consumer,当前consumer排序第二,那么就返回第三,第四个partition
10.对新的partition ownership创建临时节点并写入consumer thread id相关信息,针对partition调用ConsumerFetcherManager.startConnections,过程中会启动一个线程找zk拉取partition的leader信息,这个线程还会做一件事,addFetcherForPartitions启动ConsumerFetcherThread,ConsumerFetcherThread才是这正干活拉消息的
11.发消息使用的是SimpleConsumer,SimpleConsumer就是一个简单的blocking的rpc api封装。这里感觉kafka使用内存还是很残暴的,很多可以使用对象池和内存池的地方都没有做,new BoundedByteBufferSend,ByteBuffer.allocate(size)这样的代码随处可见。
12.KafkaStream的数据来自PartitionTopicInfo,有一个比较曲折的关联关系

同步消息

1.在ReplicaManager.makeFollowers里面启动ReplicaFetcherThread找leader拉取消息
2.通过SimpleConsumer(一个简单的rpc封装)发送FetchRequest,拿到response之后append到本地的log,然后更新一下本地的highWaterMark

创建topic

1.客户端分配partition和replica到不同的broker上面。分配方式是对partition的第一份replica按round-robin的方式开始点是随机的一个broker,第二份和后面的replica会加上一个shift,下面是源码中给的例子
broker-0  broker-1  broker-2  broker-3  broker-4
p0        p1        p2        p3        p4       (1st replica)
p5        p6        p7        p8        p9       (1st replica)
p4        p0        p1        p2        p3       (2nd replica)
p8        p9        p5        p6        p7       (2nd replica)
p3        p4        p0        p1        p2       (3nd replica)
p7        p8        p9        p5        p6       (3nd replica)

2.把topic信息和topic,partition和broker的关系写到zk上去
3.server端调用的发起点是监听了zk变更的TopicChangeListener,首先通过zk回调的数据和server缓存的数据差算出新增的topic
4.KafkaController.onNewTopicCreation,首先给topic注册AddPartitionListener,监听的路径是/brokers/topics/topicXXX,然后主动触发一次onNewPartitionCreation事件
5.进入PartitionStateMachine和ReplicaStateMachine的handleStateChange更新状态到NewPartition
6. PartitionStateMachine更新状态到OnlinePartition,在zk上初始化partition的leader和in sync replica的信息,zk path /brokers/topics/huying_test1/partitions/0/state data {"controller_epoch":6,"leader":0,"version":1,"leader_epoch":0,"isr":[0]} 这里有段注释需要注意,写zk可能会失败,当前机器可能会因为长gc失去和zk的session,所以后面catch一个zk node exist exception
7.对成为这个partition的leader的broker发送一个LeaderAndIsrRequest
8.负责管理发送的ControllerChannelManager会保存对每个broker的连接,发送也是先加到一个queue然后后面有一个线程异步发送
9. ReplicaStateMachine更新状态到OnlinePartition
10.broker收到LeaderAndIsrRequest调用ReplicaManager.becomeLeaderOrFollower,先判断controllerEpoch,如果比自己的小就丢弃这个请求,遍历请求中的每个partition,如果请求中对应partition的leaderEpoch小于当前partition的leaderEpoch就忽略当前partition
11.算出请求中partition对应的leader和follower的变更,对于变成leader的partition首先停止对这些partition的拉取job。
12.对每个partition调用Partition.makeLeader,对每个replica如果之前没有创建就创建一个,创建的过程是如果是远程的就只创建对象放在内存里,如果是local的就先创建一个Log,然后读取对应的highWatermark数据,最后生成Replica对象
13.highWatermark对应的replication-offset-checkpoint文件,和recovery point类似,但是recovery point是记录的刷盘的点,high water mark记录的是写入commit的点,所有的replica最小的commit offset就是high water mark。这里可能会增加high water mark
14.如果请求过来的topic是__consumer_offsets,那就启动OffsetManager的异步读。这个topic是用来管理所有的consumer的进度的,这样避免了把消费进度存zk上面影响扩展性。这个异步读会一直读取__consumer_offsets并把消息解码成消费进度放入缓存
15.对于需要变成follower的partition,如果是leader就调用Partition.makeFoller,首先如果本地的replica没有那就创建对应的Log,然后填充和清理一些本地内存结构
16.停止对这些成为follower的partition的拉取线程,把这些partition的Log截断到highWaterMark的位置,并启动对那些成为leader的partition的拉取线程
17.第一次创建topic的最后会启动highWaterMark的checkPoint线程,这个线程定期刷所有replica的checkPoint到磁盘


扩容

1.使用partition reassign tool
> bin/kafka-reassign-partitions.sh --zookeeper localhost:2181 --topics-to-move-json-file topics-to-move.json --broker-list "5,6" --generate 
Current partition replica assignment

{"version":1,
 "partitions":[{"topic":"foo1","partition":2,"replicas":[1,2]},
               {"topic":"foo1","partition":0,"replicas":[3,4]},
               {"topic":"foo2","partition":2,"replicas":[1,2]},
               {"topic":"foo2","partition":0,"replicas":[3,4]},
               {"topic":"foo1","partition":1,"replicas":[2,3]},
               {"topic":"foo2","partition":1,"replicas":[2,3]}]
}

Proposed partition reassignment configuration

{"version":1,
 "partitions":[{"topic":"foo1","partition":2,"replicas":[5,6]},
               {"topic":"foo1","partition":0,"replicas":[5,6]},
               {"topic":"foo2","partition":2,"replicas":[5,6]},
               {"topic":"foo2","partition":0,"replicas":[5,6]},
               {"topic":"foo1","partition":1,"replicas":[5,6]},
               {"topic":"foo2","partition":1,"replicas":[5,6]}]
}
2.PartitionsReassignedListener监听/admin/reassign_partitions,PartitionsReassignedListener会触发ReassignedPartitionsIsrChangeListener监听/brokers/topics/huying_test1/partitions/0/state
3.触发点是PartitionsReassignedListener,对每个partition先判断新分配的分区是否已经分配好,是否有对应的broker是挂掉的,这两种情况直接退出。然后注册监听器ReassignedPartitionsIsrChangeListener,再后进入onPartitionReassignment
4.KafkaController.onPartitionReassignment 初始状态新分配的replica只要有not in sync的状态areReplicasInIsr的判断就是false,首先把旧的replica和新分配的replica的并集更新一下内存和zk
5.每个partition也有一个leaderEpoch,用于ReplicaManager判断是否是旧的leaderAndISRRequest,这个属性保存在zk的原因应该是controller需要知道leaderEpoch信息
6.先更新一下zk上的leaderEpoch然后给并集所有的broker发送一个LeaderAndIsr来调整leader
7.把差集的replica设置成NewReplica的状态,这个过程也会发送LeaderAndIsr请求
8.这个过程可能会触发ReassignedPartitionsIsrChangeListener,会先判断这个partition是否结束了reassignment,然后判断是否所有reassign的replica都已经追上(isr),如果追上就再次进入onPartitionReassignment
9.把reassignedReplicas在replicaStateMachine里面转换成Online的状态,在内存里面把原来的并集替换成reassignedReplicas,如果leader不在reassignedReplicas里面通过partitionStateMachine发起一个新的partition leader选举,如果leader在的机器挂了也要重新选举,否则就只是发送请求更新一下leaderEpoch
10.把差集replica下线,下线是通过replicaStateMachine做的,依次做了几个状态变更OfflineReplica->ReplicaDeletionStarted->ReplicaDeletionSuccessful->NonExistentReplica过程中可能会停止replica fetch的线程和删掉Log
11.先把reassignedReplicas更新一下zk,然后更新一下/admin/reassign_partitions的内容,取消掉对应的ReassignedPartitionsIsrChangeListener,最后给每个broker发送一个更新metadata的请求
12.感觉这个流程主要做的是replica调整,并没有涉及到停写,和consumer消费完流程的控制,过程中可能会出现消息乱序。update,这里还真不需要,kafka一个topic多搞几个队列,扩容的时候每个机器assign的队列少一点就实现自动扩容了,缺点是文件多了,写性能可能会下降。rocketmq因为是双写而不是isr,所以扩容就需要停写和消费完。

高可用

high water mark

1.high water mark通常是一个partition所有replica的endOffset的最小值,也是同步提交模式情况下集群commit的位置
2.读的时候客户端最多只能读到highWatermark的位置
3.一个follower挂掉重启的时候首先扔掉它自己highWatermark之后的数据,然后开始追赶leader
4.leader挂掉的会重新选,新leader把自己的endOffset当做新的highWatermark,然后让其它的replica开始追赶

in sync replica

1.Partition会启动一个超时线程,调用Partition.maybeShrinkIsr,检查所有isr的end offset和最后更新时间,差距过大的就踢掉
2.处理follower的fetch request的时候可能会触发Partition.updateLeaderHWAndMaybeExpandIsr,重新检查这个replica的各种条件,可能的话就加回到isr

partition leader election

1.PartitionLeaderSelector.OfflinePartitionLeaderSelector是用得最多的leader selector,选leader的逻辑是从controller context里面拿出live isr和live replica,优先选择live isr的第一个,如果没有就选择live replica的第一个,否则就抛异常
2.更新zk path /brokers/topics/huying_test1/partitions/0/state
3.给这个partition的每个replica broker发一个请求LeaderAndIsrRequest,其中leader和follower的指令稍有区别。然后给所有活着的broker发一个请求UpdateMetadataRequest。这里的请求都是异步带callback(多数情况是null)的
4.发送请求如果失败会重试,终止重试的流程是zk那边发起,导致当前发送线程被关闭
5.接收到请求的broker开始ReplicaManager.makeLeaders, makeFollowers相关的流程
6.处理UpdateMetadataRequest请求比较简单,更新一下cache数据就ok了
7.触发点,KafkaController.onBrokerStartup, KafkaController.onBrokerFailure分别对应broker启动的时候和broker挂掉的时候,KafkaController.onPreferredReplicaElection,使用了工具
bin/kafka-preferred-replica-election.sh --zookeeper localhost:12913/kafka --path-to-json-file topicPartitionList.json

这个工具的作用是重新平衡partition leader的位置,因为broker的各种failover可能partition的leader会飘到一台机器上面去,这个工具可以使leader重新平衡。工作的原理是监听/admin/preferred_replica_election节点,触发一次partitionStateMachine.handleStateChanges,每个partition都选第一个replica当leader,由于之前分配topic,partition的时候已经在broker上面做了平衡所以这个时候leader的分配也是平衡的

controller leader election

1.controller是kafka全局的leader,controller负责监听zookeeper然后分析状态变更发送命令到其它的broker
2.broker启动的时候开始监听/controller节点,这里注意一个zk的bug(session失效可能过了一段时间临时节点才会删除,如果碰到node exist exception就返回可能会都选不上) LeaderChangeListener可能会触发KafkaController.onControllerFailover或者KafkaController.onControllerResignation
3.onControllerFailover首先会读取zk上的epoch,然后+1乐观锁方式更新zk,epoch用来标识正确的leader,当新leader产生旧leader没挂掉的时候可以作废掉旧leader的命令,broker里面关注epoch的manager保存一份epoch,接到比自己大的epoch就设定成新的
4.监听几个zk path,/admin/reassign_partitions,/admin/preferred_replica_election,/brokers/topics,/admin/delete_topics,/brokers/ids,/brokers/topics/topicXXX
5.初始化controllerContext相关的内存结构和一些manager,并给其它的broker发一个metadata更新的命令
6.如果配了自动平衡partition leader,就启动一个检查并自动平衡partition leader的scheduler,调用checkAndTriggerPartitionRebalance,不平衡率高的时候触发onPreferredReplicaElection
7.onControllerResignation执行相反的过程,取消注册之前的几个zk listener并且关闭几个manager

broker启动和关闭

1.broker启动的时候注册SessionExpireListener,这里调用的是zkClient.subscribeStateChanges(sessionExpireListener)而不是监听zk临时节点,感觉碰到机器假死会少了删临时节点的手段。然后注册临时节点/brokers/ids/xxx,这个临时节点会被controller的BrokerChangeListener监听
2.KafkaController.onBrokerStartup 首先给新启动的broker发送metadata
3.给新broker上的所有replica触发OnlineReplica的状态变更,会触发broker上面的replicaManager.becomeLeaderOrFollower,具体的逻辑和前面创建topic里面描述的相同
4.触发partitionStateMachine的状态变更,可能会导致leader重新选举,同时也会触发broker上面的replicaManager.becomeLeaderOrFollower
5.如果在新的broker上有partition reassignment,就调用onPartitionReassignment
6.broker关闭的逻辑类似,也是触发partitionStateMachine和replicaStateMachine的状态变化,稍有不同的是那些leader在这台broker上的partition要先触发一个Offline再触发一个Online

rebalance

脑裂问题和惊群效应

1.因为网络延时和不同的客户端连的zk server可能不同,不同客户端看到的zk状态可能不完全一致导致rebalance的时候的判断可能不太准确zk different client view,导致rebalance因zk ownership冲突而失败。有一个bug,[KAFKA-242]比较快的连续调用ConsumerConnector.createMessageStreams,可能会导致位点不正确前置导致丢消息,看起来发生在前一次ConsumerConnector.createMessageStreams导致rebalance结束前又调用了一次ConsumerConnector.createMessageStreams
2.一个broker或者consumer的变化可能会导致所有的consumer的rebalance。0.8版本用consumer thread id排序然后重新分配消费,中间插入一个consumer可能会导致很多consumer的分配变化,而且每次rebalance都是要关闭消费者(包括关连接和清queue),相当重。0.9 group reassignment 也是用的roundrobin 但是rebalance不会直接关闭连接

kafka 0.9的解决方案



1.增加了kaka.coordinator这个新包处理consumer的协调问题,GroupCoordinator是broker端处理所有group api请求的核心类。0.9包的客户端只有java版的
2.当一个broker变成__consumer_offsets的某个partition的leader的时候会异步的读取__consumer_offsets里面的消费位点和consume group metadata并缓存起来
3.KafkaApis里面ListGroupsKey,DescribeGroupsKey很简单就是取一下cache,HeartbeatKey也简单正常情况就是给个response然后schedule下一次heartbeat
4.SyncGroupKey,如果coordinator处于Dead,PreparingRebalance状态都是简单的不做处理,处于Stable也仅仅只是触发一下下一次心跳,处于AwaitingSync状态只处理leader client(加入group的第一个client)的请求,把当前的groupMetadata和assignment数据序列化成一条消息发到topic里面去。发送失败重新触发rebalance,发送成功本地保存一下assignment并且把状态转成stable。这里的回调看起来比较复杂,但实际上只是把返回状态码处理的逻辑传来传去
5.JoinGroupKey,加入group前先有一个简单的协议验证,协议在ConsumerProtocol比较简单,有一些版本号,topic,partition之类的信息。实际处理逻辑在GroupCoordinator.doJoinGroup进行,如果coordinator处于PreparingRebalance状态,收到请求增加或者更新一下group member然后触发下一次rebalance,处于AwaitingSync状态,如果碰到不在group里面的client那接添加member触发下一次rebalance,如果是已知的client先进行一次protocol比较也就是比较metadata,如果没有变化那就加入成功给response,否则就认为metadata变化了会更新metadata触发下一次rebalance。处于Stable状态,碰到新client或者client metadata变化都会促发rebalance。如果coordinator处于PreparingRebalance状态在最后会触发joinPurgatory.checkAndComplete把之前delay的操作刷掉。触发rebalance是往joinPurgatory里面丢一个DelayedJoin操作,DelayedJoin在tryComplete时会检查是否group里面的client都加入了,如果是会触发GroupCoordinator.onCompleteJoin,onCompleteJoin先移除失败的client,然后把generationId +1,选择一个protocol(取交集的第一个),把状态改成AwaitingSync,最后给每个client返回JoinGroupResult(有generationId,leaderId和选择的protocol)
6.LeaveGroupKey,先移除心跳配置,然后根据coordinator状态可能触发一次rebalance
7.GroupCoordinatorKey,读取group coordinator信息,首先读取topic __consumer_offsets的metadata,根据groupId算出group metadata应该在哪个partition里面,读取对应partition的partitionMetadata里面的leader信息这个leader就是group coordinator。比较重要的问题就是谁是coordinator,GroupMetadataManager.isGroupLocal->ownedPartitions.add(offsetsPartition)->loadGroupsForPartition->handleGroupImmigration->handleLeaderAndIsrRequest,所以__consumer_offsets这个topic对应的partition的leader就是group coordinator
8.KafkaConsumer每次拉消息的时候会先确保拿到server端coordinator的信息,过程中可能会发送GroupCoordinatorRequest。检查一下consumer是否在group内,如果不在就发送JoinGroupRequest,拿到response之后会发送SyncGroupRequest,前面的response会告诉consumer是不是leader,如果consumer是leader,会先调用performAssignment(这部分逻辑在WorkerCoordinator源码在外面)进行partition-consumer分配,request里面会带上group assignment,这些做完以后就实际拉数据了,比较简单
9.consumer的heartbeat启动在AbstractCoordinator.HeartbeatTask.reset,也比较简单,heartbeat失败就标记coordinator挂掉了,碰到ILLEGAL_GENERATION设置状态rejoinNeeded,这个状态位会使得消费的时候先停住发送JoinGroupRequest
10.流程串起来,consumer首先给自己知道的broker发请求得知coordinator地址,consumer通过心跳保持和coordinator的联系,如果心跳返回IllegalGeneration就意味着要rebalance,下一次消费消息会block住并且发出JoinGroupRequest,收到响应后如果是leader consumer会进行partition分配,接下来发送SyncGroupRequest,leader的request会让coordinator进入stable状态,follower会从coordinator那里拿到partition assignment,这个做完以后才开始拉新的partition
11.coordinator是__consumer_offsets这个topic的partition所在的leader机器,每次partition leader选举的时候新leader会从本地的log文件读取到GroupMetadata做缓存,consumer挂掉的时候coordinator把它移出group并且触发rebalance,coordinator通过IllegalGeneration通知consumer发生rebalance并且只有收到所有的consumer的JoinGroupRequest之后coordinator才会给consumer响应JoinGroupResponse,结束rebalance。coordinator也会监听topic partition的变化,必要的时候触发rebalance。consumer join group的时候会被coordinator分配一个consumer id之后的每次HeartbeatRequest和OffsetCommitRequest请求都带consumer id,如果哪次heartbeat请求或者commit offset请求的consumer id不对,coordinator会返回一个UnknownConsumer。

高性能

1.磁盘顺序写速度很快,pagecache的存在加速读写
2.攒消息成message set,网络上传大包,更大块的磁盘操作,内存连续性更好
3.协议固定,消息字节在broker,consumer,producer之间传输不需要修改
4.字节从pagecache到socket在linux可以通过sendfile优化

小技巧

1.SkimpyOffsetMap OffsetMap的开放式寻址实现,省内存,不能删除。
2.Throttler 一个小小的限流器,通过sleep限流。
3.多层结构的TimingWheel,解决了简单TimingWheel范围比较小的缺点。

ZK技巧

在看kafka之前让我设计broker和zk的监听结构,对于partition leader那一块我一定会这样设计

这样做的后果是当partition很多,broker也不少的时候zk上面的监听器会非常多,这样会导致zk的压力过大,而且可能zk有很多通知都是不必要的。而kafka的zk监听结构大致是这样,

broker先通过zk选出一个leader,然后只有leader监听zk上面的partition变化,监听到变化以后leader会决定是否给其它broker发送命令,这种结构zk的压力就不会增长得太厉害

ZK结构

截了一个简单的,有些需要流程触发的节点没有
/consumers
/consumers/console-consumer-84505
/consumers/console-consumer-84505/offsets
/consumers/console-consumer-84505/offsets/huying_test
/consumers/console-consumer-84505/offsets/huying_test/3
/consumers/console-consumer-84505/offsets/huying_test/2
/consumers/console-consumer-84505/offsets/huying_test/1
/consumers/console-consumer-84505/offsets/huying_test/0
/consumers/console-consumer-84505/owners
/consumers/console-consumer-84505/owners/huying_test
/consumers/console-consumer-84505/owners/huying_test/3
/consumers/console-consumer-84505/owners/huying_test/2
/consumers/console-consumer-84505/owners/huying_test/1
/consumers/console-consumer-84505/owners/huying_test/0
/consumers/console-consumer-84505/ids
/consumers/console-consumer-84505/ids/console-consumer-84505_huyingdeMacBook-Air.local-1460795322590-79d6b9c5
/config
/config/topics
/config/topics/huying_test
/config/changes
/controller
/admin
/admin/delete_topics
/brokers
/brokers/topics
/brokers/topics/huying_test
/brokers/topics/huying_test/partitions
/brokers/topics/huying_test/partitions/3
/brokers/topics/huying_test/partitions/3/state
/brokers/topics/huying_test/partitions/2
/brokers/topics/huying_test/partitions/2/state
/brokers/topics/huying_test/partitions/1
/brokers/topics/huying_test/partitions/1/state
/brokers/topics/huying_test/partitions/0
/brokers/topics/huying_test/partitions/0/state
/brokers/ids
/brokers/ids/0
/zookeeper
/zookeeper/quota
/controller_epoch

scala爽和不爽

ret.get(leaderBrokerId) match {
  case Some(element) =>
    dataPerBroker = element.asInstanceOf[HashMap[TopicAndPartition, Seq[KeyedMessage[K,Message]]]]
  case None =>
    dataPerBroker = new HashMap[TopicAndPartition, Seq[KeyedMessage[K,Message]]]
    ret.put(leaderBrokerId, dataPerBroker)
}
模式匹配,在接收请求的时候用起来特别爽
brokerIds.map { brokerId =>
  partitionReplicaAssignment
    .filter { case(topicAndPartition, replicas) => replicas.contains(brokerId) }
    .map { case(topicAndPartition, replicas) =>
             new PartitionAndReplica(topicAndPartition.topic, topicAndPartition.partition, brokerId) }
}.flatten.toSet
server端编程大量集合操作的时候这种高阶流式处理集合比java方便太多
val threadAssignor = Utils.circularIterator(headThreadIdSet.toSeq.sorted)
这里用到了Stream的lazy evaluation,制造了一个无限循环的队列,cool!

不爽的地方
1.debug的时候ide支持还是不够,有些表达进不去或者看不到值
2.有些写法写起来爽读起来痛苦,比如这样的
val streams = e._2.map(_._2._2).toList

发送端性能问题

Sarama假异步

sarama客户端本身是支持异步发送模型的,但是sarama老版本实际上做的假异步https://github.com/Shopify/sarama/issues/2103
在后续的版本中做了修复,https://github.com/Shopify/sarama/releases/tag/v1.31.0

服务端请求处理

1.网络层将请求投递到 RequestChannel 
2.多线程的 KafkaRequestHandler 并发的从 RequestChannel 获取请求, 交给 kafkaApis 处理

这里看起来会有顺序问题,但是processCompletedReceives里面会调用selector.mute禁读,所以服务端对同一个连接上的读请求是一个一个顺序处理的过程。这个时候在tcp buffer里面是可以缓存更多的请求的,通过这个buffer可以平摊掉客户端到服务端的网络延迟

WaitForAll

因为服务端对同一个连接上的请求是顺序处理的,所以WaitForAll的时候broker间数据复制的延迟会让服务端处理一个请求的耗时变长而这个耗时是没法通过异步的方式化解的。batch还是有可能平摊掉多个请求顺序处理的时间,但如果slave broker的io有问题batch效果可能也不明显?