【问题标题】:Checkpoint with spark file streaming in javajava中带有spark文件流的检查点
【发布时间】:2021-01-07 15:46:24
【问题描述】:

如果在任何情况下我的 spark 流应用程序停止/终止,我想使用 spark 文件流应用程序实现检查点以处理来自 hadoop 的所有未处理文件。我正在关注这个:streaming programming guide,但没有找到 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();
        }   
    }
}

【问题讨论】:

    标签: java hadoop spark-streaming


    【解决方案1】:

    你必须使用Checkpointing

    对于检查点,使用updateStateByKeyreduceByKeyAndWindow有状态 转换。 spark-examples 中有很多示例,以及 git-hub 中的预构建 spark 和 spark 源。具体见JavaStatefulNetworkWordCount.java;

    【讨论】:

      猜你喜欢
      • 2023-02-09
      • 2016-04-05
      • 1970-01-01
      • 1970-01-01
      • 2021-06-02
      • 2015-06-16
      • 1970-01-01
      • 2017-08-23
      • 1970-01-01
      相关资源
      最近更新 更多