SparkStreaming如何处理消费Kafka的数据以确保无丢失无重复
发布时间
阅读量:
阅读量
目录
SparkStreaming对Kafka数据的接收有两种实现方式:]()Receiver用于接收数据,并采用了Direct 方式。
某一个接收器运行效率较低,在处理大量消息时必须启动多个线程进行处理,并手动合并收集到的数据后再进行后续处理。为了确保无数据丢失这一目标,WAL(预写日志)机制被采用以保证系统的安全性。通过此机制,系统会将所有接收到的Kafka消息同步存储至HDFS,以便在出现故障时恢复所有存储的数据。值得注意的是,WAL虽然能够实现无遗漏的目标,但却无法保证'恰好一次'的完整性。当接收器完成消息收集并将它们保存至HDFS后,在offset值还未被更新至ZooKeeper之前,系统会认为消息未被正确接收。一旦接收器重新启动,系统会从WAL中读取备份并开始消费这些备份的消息,因为此时系统会基于ZooKeeper中最新的offset值来进行消费操作。
(2)该Direct方法依赖于checkpoint机制来确保其稳定性。每当Spark Streaming从Kafka消费数据时,系统会将所消费的Kafka偏移量同步至checkpoint存储,从而消除与ZooKeeper不一致的问题,并且可以直接从各个分区读取数据,显著提升了并行处理的能力。当程序发生故障或需要升级时,系统便可以继续上次的数据读取,从而确保数据在断线时不会丢失。然而Checkpoi
全部评论 (0)
还没有任何评论哟~
