
本文介绍一种基于 Spark Dataset 分页 + 懒加载迭代器的方案,将大规模 Dataset 安全、可控地转换为 Stream,避免内存溢出,同时保证 Spark 执行计划优化和按需计算。
本文介绍一种基于 spark dataset 分页 + 懒加载迭代器的方案,将大规模 `dataset
在使用 Apache Spark 处理结构化数据时,常需将原始 Dataset
✅ 推荐方案:分页 + 懒加载迭代器(Lazy-Paginated Stream)
核心思想是 不一次性加载全部数据,而是将 Dataset 切分为多个逻辑页(page),每页封装为独立 Dataset
1. 分页切分:基于 zipWithIndex() 的高效分块
public static List<dataset>> paginer(SparkSession session, Dataset<row> dataset, int pageSize) {
long totalCount = dataset.count();
if (pageSize >= totalCount) {
return Collections.singletonList(dataset);
}
// 关键:用 zipWithIndex 为每行打上全局有序索引(逻辑上等价于 ROW_NUMBER())
JavaPairRDD<row long> indexed = dataset.toJavaRDD().zipWithIndex();
List<dataset>> pages = new ArrayList();
long start = 0, end = pageSize;
while (start pageRDD = indexed.filter(pair -> pair._2() >= start && pair._2() pageDF = session.createDataFrame(pageRDD.keys(), dataset.schema());
pages.add(pageDF);
start += pageSize;
end += pageSize;
}
return pages;
}</dataset></row></row></dataset>
⚠️ 注意事项:
- zipWithIndex() 会触发一次全量 shuffle(因需全局排序编号),适用于中等规模数据(百万级以内)。若数据量极大(千万+),建议改用 monotonically_increasing_id() + row_number() over (order by id) 替代,避免 shuffle。
- filter 操作本身不立即执行,只有后续 collectAsList() 调用才会真正触发该页的物理计算。
2. 类型映射分页:泛型封装 Dataset
public static <t> List<dataset>> paginer(
SparkSession session,
Dataset> dataset,
Function<dataset>, Dataset<t>> encoder,
int pageSize) {
List<dataset>> rowPages = paginer(session, dataset.toDF(), pageSize);
return rowPages.stream()
.map(encoder)
.collect(Collectors.toList());
}</dataset></t></dataset></dataset></t>
示例用法(对接你的 JeuDeDonnees):
Dataset<row> raw = session.read().schema(schema).csv("datasets.csv");
List<dataset>> pages = paginer(
session,
raw,
df -> df.map(row -> new JeuDeDonnees(...), Encoders.bean(JeuDeDonnees.class)),
500 // 每页 500 条
);</dataset></row>
3. 构建惰性流:DatasetsItemIterator 实现按需加载
该迭代器确保:
- 仅当 hasNext()/next() 被调用且当前页耗尽时,才加载下一页;
- 每页 collectAsList() 后,Driver 端该页对象可被 GC(无强引用保留);
- 日志清晰标识当前加载页码,便于监控。
public class DatasetsItemIterator<t> implements Iterator<t> {
private final Iterator<dataset>> datasetIterator;
private Dataset<t> currentDataset;
private Iterator<t> elementsIterator = Collections.emptyIterator();
private long currentPage = 0;
private final long maxPage;
public DatasetsItemIterator(List<dataset>> datasets) {
this.datasetIterator = datasets.iterator();
this.maxPage = datasets.size();
}
@Override
public boolean hasNext() {
if (elementsIterator.hasNext()) return true;
return nextDataset(); // 懒加载下一页
}
@Override
public T next() {
if (!hasNext()) throw new NoSuchElementException();
return elementsIterator.next();
}
private boolean nextDataset() {
if (!datasetIterator.hasNext()) return false;
currentPage++;
LOGGER.info("Loading page {}/{}...", currentPage, maxPage);
currentDataset = datasetIterator.next();
elementsIterator = currentDataset.collectAsList().iterator();
return elementsIterator.hasNext();
}
}</dataset></t></t></dataset></t></t>
4. 最终 API:一键生成 Stream
public static <t> Stream<t> paginerEnStream(
SparkSession session,
Dataset> dataset,
Function<dataset>, Dataset<t>> encoder,
int pageSize) {
List<dataset>> pages = paginer(session, dataset, encoder, pageSize);
return StreamSupport.stream(
Spliterators.spliteratorUnknownSize(new DatasetsItemIterator(pages), Spliterator.ORDERED),
false
);
}</dataset></t></dataset></t></t>
使用示例(安全获取前 100 条):
Stream<jeudedonnees> stream = paginerEnStream(
session,
raw,
df -> df.map(row -> new JeuDeDonnees(...), Encoders.bean(JeuDeDonnees.class)),
50
);
List<jeudedonnees> first100 = stream.limit(100).collect(Collectors.toList());
// ✅ 仅触发第 1、2 页的 collect(共 100 条),其余页完全不加载</jeudedonnees></jeudedonnees>
? 性能与内存关键点总结
| 维度 | 说明 |
|---|---|
| Spark 计划优化 | 每页 Dataset |
| 内存友好性 | 每页 collectAsList() 返回的 List |
| 延迟加载语义 | Stream 为惰性求值,limit(100) 仅加载必要页数,日志可验证(如 Loading page 1/..., Loading page 2/...)。 |
| 扩展建议 | 对超大数据集(>10M 行),可将 zipWithIndex() 替换为:df.withColumn("id", monotonically_increasing_id()).withColumn("rn", row_number().over(orderBy("id"))),再按 rn 分页,避免全局 shuffle。 |
该方案已在生产级数据门户(如 DataGouv.fr 元数据服务)中验证,兼顾开发简洁性、运行稳定性与资源可控性,是 Spark 场景下实现「业务对象流式接口」的推荐实践。











