
# 配置Spark Streaming对接Kafka时数据后进先出功能
#### 操作场景
当Spark Streaming应用与Kafka对接，Spark Streaming应用异常终止并从checkpoint恢复重启后，对于进入Kafka数据的任务，系统默认优先处理应用终止前（A段时间）未完成的任务和应用终止到重启完成这段时间内（B段时间）进入Kafka数据生成的任务，最后再处理应用重启完成后（C段时间）进入Kafka数据生成的任务。并且对于B段时间进入Kafka的数据，Spark将按照终止时间（batch时间）生成相应个数的任务，其中第一个任务读取全部数据，其余任务可能不读取数据，造成任务处理压力不均匀。
如果A段时间的任务和B段时间任务处理得较慢，则会影响C段时间任务的处理。针对上述场景，Spark提供Kafka后进先出功能。
图1Spark Streaming应用重启时间轴   
![](https://support.huaweicloud.com/cmpntguide-lts-mrs/zh-cn_image_0000001676748598.png "点击放大")
开启此功能后，Spark将优先调度C段时间内的任务，如果存在多个C段任务，则按照任务产生的先后顺序调度执行，再执行A段时间和B段时间的任务。另外，对于B段时间进入Kafka的数据，Spark除了按照终止时间生成相应任务，还将这个期间进入Kafka的所有数据均匀分配到各个任务，避免任务处理压力不均匀。
约束条件：
- 目前该功能只适用于Spark Streaming中的Direct方式，且执行结果与上一个batch时间处理结果没有依赖关系（即无state操作，如updatestatebykey）。对多条数据输入流，需要相对独立无依赖的状态，否则可能导致数据切分后结果发生变化。
- Kafka后进先出功能的开启要求应用只能对接Kafka输入源。
- 如果提交应用的同时开启Kafka后进先出和流控功能，对于B段时间进入Kafka的数据，将不启动流控功能，以确保读取这些数据的任务调度优先级最低。应用重新启动后C段时间的任务启用流控功能。
 
#### 配置描述
1. 安装Spark客户端。 
   详细操作请参考[安装MRS客户端](https://support.huaweicloud.com/cmpntguide-lts-mrs/mrs_01_2127.html)。
   
   
2. 使用客户端安装用户登录Spark客户端节点。 
   在Spark客户端的"*{客户端安装目录}*/Spark/spark/conf/spark-defaults.conf"配置文件中进行设置，修改如下参数：
   表1参数说明 
   | 参数                                         | 参数说明                                                                                                                                                                                                                                                                                                                                                                    | 取值示例                                                          |
   |:---|:---|:---|
   | spark.streaming.kafka.direct.lifo          | Spark Streaming中控制Kafka的数据处理顺序。该参数通过调整数据读取策略，影响实时流处理的延迟和吞吐量。 - false：优先处理最早接收到的 Kafka 消息。适合需要严格按时间顺序处理数据的场景，但可能增加最新数据的处理延迟。  - true：优先处理最新接收到的Kafka消息。大幅降低最新数据的处理延迟，但可能忽略历史未处理数据。   | true                                                          |
   | spark.streaming.kafka010.inputstream.class | 获取解耦在FusionInsight侧的类。                                                                                                                                                                                                                                                                                                                                                  | org.apache.spark.streaming.kafka010.xxDirectKafkaInputDStream |
      
   
   
 
