SparkStreaming的Receive和Direct模式以及背压机制

时间:2024-05-21 18:58:48

Receive和Direct模式对比

Receiver模式
SparkStreaming的Receive和Direct模式以及背压机制SparkStreaming的Receive和Direct模式以及背压机制

  • 数据是源源不断的通过 receiver 接收,当数据被接收后,其将这些数据存储在 Block Manager 中;为了不丢失数据,其还将数据备份到其他的 Block Manager 中;
  • Receiver Tracker 收到被存储的 Block IDs,然后其内部会维护一个时间到这些 block IDs 的关系;
  • Job Generator 会每隔 batchInterval 的时间收到一个事件,其会根据这段时间到来的数据和stage生成一个 JobSet;
  • Job Scheduler 运行上面生成的 JobSet,将JobSet分发到对应的executor上运行。
  • 存在的问题:当batch processing time>batchinterval 这种情况持续过长的时间,会造成数据在内存中堆积,导致Receiver所在Executor内存溢出等问题;

Direct模式
SparkStreaming的Receive和Direct模式以及背压机制

  • 与receiver模式类似,不同在于没有单独receiver组件,各个executor根据自己消费速率直接从kafka中拉取数据,这样避免了SparkStreaming与数据源生产速率不均衡造成的数据积压。
  • 同时可以自己维护kafka的offset,避免数据丢失

Receiver的背压机制

为什么需要背压机制?
当batch processing time>batchinterval 这种情况持续过长的时间,会造成数据在内存中堆积,导致Receiver所在Executor内存溢出等问题;

1.5之后的背压体系结构
SparkStreaming的Receive和Direct模式以及背压机制

  • 新增了一个RateController实现自动调节数据的传输速率。基于 processingDelay 、schedulingDelay 、当前 Batch 处理的记录条数以及处理完成事件来估算出一个速率;这个速率主要用于更新流每秒能够处理的最大记录的条数。
  • InputDStreams 内部的 RateController 里面会存下计算好的最大速率,这个速率会在处理完 onBatchCompleted 事件之后将计算好的速率推送到 ReceiverSupervisorImpl,这样接收器就知道下一步应该接收多少数据了。