Apache Beam Spark runner的side inputs引起了SIGNAL TERM。
创始人
2024-09-03 13:31:10
0

Apache Beam的Spark runner在使用side inputs时可能会导致SIGNAL TERM错误。这个问题通常是由于Spark的任务超时时间过短导致的。为了解决这个问题,你可以通过增加Spark任务的超时时间来解决。以下是一个示例代码:

import org.apache.beam.runners.spark.SparkPipelineOptions;
import org.apache.beam.runners.spark.SparkRunner;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.Optional;
import org.apache.spark.api.java.function.Function2;

public class SideInputExample {
    public static void main(String[] args) {
        SparkConf sparkConf = new SparkConf().setAppName("SideInputExample");
        JavaSparkContext jsc = new JavaSparkContext(sparkConf);
        
        SparkPipelineOptions options = PipelineOptionsFactory.as(SparkPipelineOptions.class);
        options.setRunner(SparkRunner.class);
        options.setSparkContext(jsc);
        
        Pipeline pipeline = Pipeline.create(options);
        
        // 读取主输入数据
        pipeline.apply(TextIO.read().from("input.txt"))
                // 处理主输入数据
                .apply(ParDo.of(new DoFn() {
                    @ProcessElement
                    public void processElement(ProcessContext c) {
                        // 从side input中获取数据
                        Optional sideInputValue = c.sideInput("sideInput");
                        
                        // 处理主输入数据和side input数据
                        // ...
                        
                        c.output("output");
                    }
                }).withSideInputs(jsc.emptyRDD()))
                // 输出结果数据
                .apply(TextIO.write().to("output.txt").withoutSharding());
        
        // 创建side input数据
        JavaPairRDD sideInput = jsc.parallelizePairs(Arrays.asList(
                new Tuple2<>("key1", "value1"),
                new Tuple2<>("key2", "value2")
        ));
        
        // 将side input数据作为Broadcast变量
        jsc.broadcast(sideInput).toJavaRDD()
                .mapToPair(pair -> pair)
                .rdd()
                .toJavaRDD()
                .mapToPair(pair -> pair);
        
        // 运行Beam Pipeline
        pipeline.run();
    }
}

在上面的代码中,我们通过设置options中的SparkPipelineOptions来配置Spark runner。然后,我们使用pipeline.apply()来读取和处理主输入数据,并使用ParDo.of()来处理主输入数据和side input数据。在这个例子中,我们使用了一个空的RDD作为side input,你可以根据实际情况修改成你的side input数据。注意,我们使用了jsc.broadcast()方法将side input数据作为Broadcast变量,这样可以在Spark任务中共享这些数据。

最后,我们调用pipeline.run()来运行Beam pipeline。这样,Spark runner就可以正确处理side inputs,而不会导致SIGNAL TERM错误。

相关内容

热门资讯

一分钟了解!边锋老友二打一有挂... 一分钟了解!边锋老友二打一有挂,wepoker私人局透视方法,详细有挂(有挂攻略)-哔哩哔哩是一款可...
实测分享!边锋麻将有挂(weP... 实测分享!边锋麻将有挂(wePOke),太坑了的确是真的有挂(有挂方法)-哔哩哔哩是一款可以让一直输...
一秒答解!广东雀神挂件去哪买(... 相信很多朋友都在电脑上玩过广东雀神挂件去哪买吧,但是很多朋友都在抱怨用电脑玩起来不方便。为此小编给大...
必备攻略(wpk一直输)外挂透... 1、必备攻略(wpk一直输)外挂透明挂辅助APP(线上)德州ai机器人(2025已更新)(哔哩哔哩)...
分享个大家!边锋游戏辅助器,x... 分享个大家!边锋游戏辅助器,xpoker辅助器,详细有挂(有挂总结)-哔哩哔哩1、点击下载安装,微扑...
今日百科!川麻圈辅助器手机版(... 今日百科!川麻圈辅助器手机版(wePoke),太坑了其实真的有挂(有挂介绍)-哔哩哔哩;值得一提的是...
分享个大家!雀神辅助器苹果版靠... 分享个大家!雀神辅助器苹果版靠谱(辅助挂)太坑了果真是真的有挂(有挂攻略)-哔哩哔哩;人气非常高,a...
新手必备(aapOker)外挂... 新手必备(aapOker)外挂透明挂辅助机制(ai代打)发牌规律(2021已更新)(哔哩哔哩);实战...
一分钟揭秘!边锋杭麻圈辅助,h... 一分钟揭秘!边锋杭麻圈辅助,hhpoker德州有挂,详细有挂(有挂教学)-哔哩哔哩 科技详细教程;7...
分享一款!边锋麻将有挂(WeP... 《分享一款!边锋麻将有挂(WePoKer),太坑了确实是真的有挂(有挂攻略)-哔哩哔哩》 边锋麻将有...