【问题标题】:SparkSession null point exception in Dataset foreach数据集 foreach 中的 SparkSession 空点异常
【发布时间】:2021-09-23 23:45:19
【问题描述】:

我是 Spark 的新手。

我想继续从kafka获取消息,然后在消息大小超过100000时保存到S3。

我通过 Dataset.collectAsList() 实现了它,但它使用 Total size of serialized results of 3 tasks (1389.3 MiB) is bigger than spark.driver.maxResultSize 抛出错误

于是转而使用foreach,使用SparkSession创建DataFrame时出现空点异常。

你知道吗?谢谢。

---代码---

SparkSession spark = generateSparkSession();
        registerUdf4AddPartition(spark);
        Dataset<Row> dataset = spark.readStream().format("kafka")
                .option("kafka.bootstrap.servers", args[0])
                .option("kafka.group.id", args[1])
                .option("subscribe", args[2])
                .option("kafka.security.protocol", SecurityProtocol.SASL_PLAINTEXT.name)
                .load();
        DataStreamWriter<Row> console = dataset.toDF().writeStream().foreachBatch((rawDataset, time) -> {
            Dataset<Row> rowDataset = rawDataset.selectExpr("CAST(value AS STRING)");
            //using foreach
            rowDataset.foreach(row -> {
                List<Span> rawDataList = new CsvToBeanBuilder(new StringReader(row.getString(0))).withType(Span.class).build().parse();
                spans.addAll(rawDataList);
                batchSave(spark);
            });

            // using collectAsList
            List<Row> rows = rowDataset.collectAsList();
            for (Row row : rows) {
                List<Span> rawDataList = new CsvToBeanBuilder(new StringReader(row.getString(0))).withType(Span.class).build().parse();
                spans.addAll(rawDataList);
                batchSave(spark);
            }
        });
        StreamingQuery start = console.start();
        start.awaitTermination();

public static void batchSave(SparkSession spark){
        synchronized (spans){
            if(spans.size() == 100000){
                System.out.println(spans.isEmpty());
                Dataset<Row> spanDataSet = spark.createDataFrame(spans, Span.class);
                Dataset<Row> finalResult = addCustomizedTimeByUdf(spanDataSet);

                StringBuilder pathBuilder = new StringBuilder("s3a://fwk-dataplatform-np/datalake/log/FWK/ART2/test/leftAndRight");
                finalResult.repartition(1).write().partitionBy("year","month","day","hour").format("csv").mode("append").save(pathBuilder.toString());
                spans.clear();
            }
        }
    }

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-streaming spark-structured-streaming


    【解决方案1】:

    由于主SparkSession驱动程序中运行,而foreach...中的任务在executors中运行,所以spark没有定义给所有其他执行者。

    顺便说一句,在foreach 任务中使用synchronized 没有任何意义,因为一切都是分布式的。

    【讨论】:

    • 所以在foreah中时无法使用spark创建dataframe?你能告诉我对我的要求有什么想法吗? consume message from kafka and save once message size over 100000
    • 1.您可以像在驱动程序中一样在执行程序中定义任何内容(包括SparkSession); 2. 为什么你需要每 100000 条消息溢出/切片数据?
    • 如我所说,我的 spark job 会继续消费消息,如果我不控制文件大小,文件会太大
    • 你可以看看这个官方文档spark.apache.org/docs/latest/…;也许您需要检查是否可以在 s3 服务中附加数据
    猜你喜欢
    • 1970-01-01
    • 2017-05-30
    • 2011-10-19
    • 2022-07-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-08-09
    相关资源
    最近更新 更多