spark streaming 的 执行流程 分析
发布时间
阅读量:
阅读量
文章目录
-
简介
- 部署流计算引擎
- 通过调用StreamingContext.start来实现流计算的初始化
- 通过调用JobScheduler.start来配置并运行Job调度器
-
接收到与存储的数据
- 在Driver端创建一个
ReceiverTracker实例- 将Driver端的
Selector打包成一个RDD,并触发Spark作业以启动receiver - 在Executor端开始处理来自Selector的数据,并将其整理后保存
- 将Driver端的
- Blocking组件创建一个BlockingQueue来组织将要到达的block块,并完成收集与处理工作。
- 在Driver端创建一个
-
生产Batch Job用于处理数据
-
按照batchDuration的时间间隔生产相应的GenerateJobs消息
-
生产相应的GenerateJobs消息并将其提交至队列中
简介
SparkStreaming为每个数据源独立分配对应的接收器(Receiver),这些接收器作为任务形式运行于应用的Executor进程中,并通过从输入端接收数据后将其按小批量(batch)分组并存储为RDD(Resilient Distributed Datasets),最后将作业提交给Spark作业执行系统进行处理。流计算的执行流程图如图所示,并主要包含以下
全部评论 (0)
还没有任何评论哟~
