(版本定制)第11课:Spark Streaming源码解读
admin
2023-01-31 11:48:21
0

本期内容:

    1、ReceiverTracker的架构设计

    2、消息循环系统

    3、ReceiverTracker具体实现


上节课讲到了Receiver是如何不断的接收数据的,并且接收到的数据的元数据会汇报给ReceiverTracker,下面我们看看ReceiverTracker具体的功能及实现。

ReceiverTracker主要的功能:

  1. 在Executor上启动Receivers。

  2. 停止Receivers 。

  3. 更新Receiver接收数据的速度(也就是限流)

  4. 不断的等待Receivers的运行状态,只要Receivers停止运行,就重新启动Receiver,也就是Receiver的容错功能。

  5. 接受Receiver的注册。

  6. 借助ReceivedBlockTracker来管理Receiver接收数据的元数据。

  7. 汇报Receiver发送过来的错误信息


ReceiverTracker 管理了一个消息通讯体ReceiverTrackerEndpoint,用来与Receiver或者ReceiverTracker 进行消息通信。

在ReceiverTracker的start方法中,实例化了ReceiverTrackerEndpoint,并且在Executor上启动Receivers。

启动Receivr,其实是ReceiverTracker给ReceiverTrackerEndpoint发送了一个本地消息,ReceiverTrackerEndpoint将Receiver封装成RDD以job的方式提交给集群运行。

Receiver启动后,会向ReceiverTracker注册,注册成功才算正式启动了。

当Receiver端接收到数据,达到一定的条件需要将数据写入BlockManager,并且将数据的元数据汇报给ReceiverTracker。

/** Store block and report it to driver */
def pushAndReportBlock(
    receivedBlock: ReceivedBlock,
metadataOption: Option[Any],
blockIdOption: Option[StreamBlockId]
  ) {
val blockId = blockIdOption.getOrElse(nextBlockId)
val time = System.currentTimeMillis
val blockStoreResult = receivedBlockHandler.storeBlock(blockId, receivedBlock)
  logDebug(s"Pushed block $blockId in ${(System.currentTimeMillis - time)} ms")
val numRecords = blockStoreResult.numRecords
val blockInfo = ReceivedBlockInfo(streamId, numRecords, metadataOption, blockStoreResult)
trackerEndpoint.askWithRetry[Boolean](AddBlock(blockInfo))
  logDebug(s"Reported block $blockId")
}

当ReceiverTracker收到元数据后,会在线程池中启动一个线程来写数据

case AddBlock(receivedBlockInfo) =>
if (WriteAheadLogUtils.isBatchingEnabled(ssc.conf, isDriver = true)) {
walBatchingThreadPool.execute(new Runnable {
override def run(): Unit = Utils.tryLogNonFatalError {
if (active) {
          context.reply(addBlock(receivedBlockInfo)) 
        } else {
throw new IllegalStateException("ReceiverTracker RpcEndpoint shut down.")
        }
      }
    })
  } else {
    context.reply(addBlock(receivedBlockInfo))
  }

数据的元数据是交由ReceivedBlockTracker管理的

数据最终被写入到streamIdToUnallocatedBlockQueues中,一个流对应一个数据块信息的队列。

每当Streaming 触发job时,会将队列中的数据分配成一个batch,并将数据写入timeToAllocatedBlocks数据结构。

下面是简单的流程图:

(版本定制)第11课:Spark Streaming源码解读

相关内容

热门资讯

我们的“文化体力”是如何被耗尽... 《花束般的恋爱》2026年已经过半,要盘点前半年爆火的流行语,“赛博确诊”一定是其中之一。大批年轻人...
跨国药企正在疯狂买中国创新药,... ► 文 观察者网心智观察所2026年的中国创新药行业,正在发生一个耐人寻味的场景。一边,全球最大的跨...
知名装机up主被曝负债200万... 据都市快报消息:7月21日,深耕高端PC圈的B站UP主(B站、小红书、抖音等平台上传视频的创作者、博...
AI脸让人生理性厌恶?“恐怖谷... 【文/观察者网 夏曼波】最近,“对AI脸感到不适”相关话题频繁上热搜,网友纷纷表示“不适”、“生理性...
要管管国会山股神了?美国参议院... 作者 | 第一财经 樊志菁当地时间周三,美国议会众议院以232票对198票表决结果,通过一项名为《停...
济南一广场突现“知了大军”密密... 近日,有济南市民在社交媒体发布视频称:在当地彩石附近的一处广场地面上密密麻麻聚集了很多类似知了的昆虫...
伊朗军队:打击科威特境内多处美... 当地时间23日,伊朗军队在一份公告中指出,几小时前,伊朗军队对科威特多哈营地的美军弹药和后勤仓库、阿...
洛阳“珠宝大王”年永安沉浮:涉... 澎湃新闻记者 王鑫 实习生 林晨 林宜之河南洛阳的“珠宝大王”年永安淡出公众视野多年后,这个曾把“金...
聚焦核心AI项目,亚马逊裁撤其... 亚马逊在本周三宣布新动态,公司正在裁撤通用人工智能(AGI)部门的部分员工,这标志着亚马逊在持续向人...
上海交大发布第七届“十大科技进... 7月22日,上海交通大学第七届“十大科技进展”在闵行大零号湾成果转化中心发布。本届入选成果包括6项基...