Checkpoint พร้อมสตรีมไฟล์ spark ใน java

Sep 21 2020

ฉันต้องการใช้จุดตรวจด้วยแอปพลิเคชั่นสตรีมไฟล์ spark เพื่อประมวลผลไฟล์ที่ยังไม่ได้ประมวลผลทั้งหมดจาก hadoop หากในกรณีใด ๆ แอปพลิเคชันสตรีมมิ่ง Spark ของฉันจะหยุด / ยุติ ฉันกำลังทำตามสิ่งนี้: คู่มือการเขียนโปรแกรมสตรีมมิ่งแต่ไม่พบ JavaStreamingContextFactory โปรดช่วยฉันฉันควรทำอย่างไร

รหัสของฉันคือ

public class StartAppWithCheckPoint {

    public static void main(String[] args) {
        
        try {
            
            String filePath = "hdfs://Master:9000/mmi_traffic/listenerTransaction/2020/*/*/*/"; 
            String checkpointDirectory = "hdfs://Mongo1:9000/probeAnalysis/checkpoint";
            SparkSession sparkSession = JavaSparkSessionSingleton.getInstance();

            JavaStreamingContextFactory contextFactory = new JavaStreamingContextFactory() {
                  @Override public JavaStreamingContext create() {
                      
                    SparkConf sparkConf = new SparkConf().setAppName("ProbeAnalysis");
                    JavaSparkContext sc = new JavaSparkContext(sparkConf);  
                    JavaStreamingContext jssc = new JavaStreamingContext(sc, Durations.seconds(300));
                    JavaDStream<String> lines = jssc.textFileStream(filePath).cache();
                    
                    jssc.checkpoint(checkpointDirectory);
                    return jssc;
                  }
                };
                
            JavaStreamingContext context = JavaStreamingContext.getOrCreate(checkpointDirectory, contextFactory);
            
            context.start();
            context.awaitTermination();
            context.close();
            sparkSession.close();
            
        } catch(Exception e) {
            e.printStackTrace();
        }   
    }
}

คำตอบ

1 majidhajibaba Sep 22 2020 at 00:01

คุณต้องใช้Checkpointing

สำหรับ checkpointing ใช้statefulการเปลี่ยนแปลงอย่างใดอย่างหนึ่งหรือupdateStateByKey reduceByKeyAndWindowมีตัวอย่างมากมายในตัวอย่างประกายไฟที่ให้มาพร้อมกับการสร้างประกายไฟล่วงหน้าและแหล่งกำเนิดประกายไฟใน git-hub สำหรับข้อมูลเฉพาะของคุณโปรดดูที่JavaStatefulNetworkWordCount.java ;