You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java环境下Spark如何过滤Dataset列表值并拼接指定内容?

Java实现方案:过滤并拼接NewsArticle的段落内容

实体类定义(与你的需求一致)

public class NewsArticle {
    private String id;
    private String title;
    private List<ContentItem> contents;

    // 请按需添加构造器、getter和setter方法
}

public class ContentItem {
    private String content;
    private String subtype;
    private String url;

    // 请按需添加构造器、getter和setter方法
}

基于Spark Dataset的处理方案

如果你的NewsArticle Dataset是Spark环境下的数据集,可以用以下方式处理,输出指定结构的Row数据集:

import org.apache.spark.api.java.function.MapFunction;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructType;

import java.util.Collections;
import java.util.Optional;
import java.util.stream.Collectors;

public class NewsArticleProcessor {

    public Dataset<Row> processNewsArticles(Dataset<NewsArticle> newsArticleDataset) {
        // 定义输出数据的Schema结构
        StructType outputSchema = new StructType()
                .add("id", DataTypes.StringType)
                .add("title", DataTypes.StringType)
                .add("mergedContent", DataTypes.StringType);

        return newsArticleDataset.map((MapFunction<NewsArticle, Row>) article -> {
            // 过滤subtype为"paragraph"的内容,拼接成字符串,同时处理contents为null的情况
            String mergedContent = Optional.ofNullable(article.getContents())
                    .orElse(Collections.emptyList())
                    .stream()
                    .filter(item -> "paragraph".equals(item.getSubtype()))
                    .map(ContentItem::getContent)
                    .collect(Collectors.joining("\n")); // 用换行分隔段落,可替换为空格等其他分隔符

            return RowFactory.create(article.getId(), article.getTitle(), mergedContent);
        }, outputSchema);
    }
}

普通Java集合处理方案

如果只是处理普通Java集合中的NewsArticle对象,逻辑如下:

import java.util.Collections;
import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;

public class NewsArticleHandler {

    // 处理单个文章,返回拼接后的段落内容
    private String mergeParagraphs(NewsArticle article) {
        return Optional.ofNullable(article.getContents())
                .orElse(Collections.emptyList())
                .stream()
                .filter(item -> "paragraph".equals(item.getSubtype()))
                .map(ContentItem::getContent)
                .collect(Collectors.joining(" "));
    }

    // 批量处理,返回[id, title, 拼接内容]的列表
    public List<Object[]> batchProcessArticles(List<NewsArticle> articleList) {
        return articleList.stream()
                .map(article -> new Object[]{
                        article.getId(),
                        article.getTitle(),
                        mergeParagraphs(article)
                })
                .collect(Collectors.toList());
    }
}

注意事项

  • 增加了Optional处理contents为null的情况,避免空指针异常
  • 拼接时的分隔符可根据业务需求调整(比如换行、空格、分号等)
  • Spark环境下需确保NewsArticle和ContentItem已正确注册为Spark可序列化的Bean,或使用对应的Encoder

内容的提问来源于stack exchange,提问作者Nico

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 08:59:19