设计模式详解-迭代器模式
设计模式详解:迭代器模式
一、模式概述
迭代器模式(Iterator Pattern)是行为型设计模式中最具遍历抽象价值的模式,其核心意图在于提供一种方法顺序访问一个聚合对象中的各个元素,而又不需要暴露该对象的内部表示。这一模式将遍历行为从聚合对象中分离出来,封装为独立的迭代器对象,使客户端能够以统一的方式遍历各种不同的数据结构。
迭代器模式的命名直接揭示了其功能——"迭代"即重复执行、逐步推进的过程。在数学中,迭代法通过重复应用函数逼近解;在软件开发中,迭代器通过重复调用next()方法遍历集合。这一模式的深层价值在于封装遍历的复杂性:无论是数组、链表、树、图,还是数据库游标、流式数据、分布式分页结果,迭代器都为客户端提供一致的hasNext()/next()接口,隐藏底层数据结构的差异与访问细节。
迭代器模式是现代编程语言的基石设施。Java的Iterator接口、C#的IEnumerable、Python的__iter__、JavaScript的Symbol.iterator,均是迭代器模式的标准化实现。这些语言级支持使迭代器模式从显式设计模式演变为隐式编程习惯,但其设计思想依然是理解集合框架、流式API、响应式编程的关键。
二、模式结构
迭代器模式包含两个核心角色,形成遍历与聚合的分离:
抽象迭代器(Iterator):定义访问和遍历元素的接口,通常包括hasNext()、next()、remove()等方法。
具体迭代器(Concrete Iterator):实现迭代器接口,维护遍历状态,跟踪聚合中的当前位置。
抽象聚合(Aggregate):定义创建迭代器对象的接口,如iterator()方法。
具体聚合(Concrete Aggregate):实现聚合接口,返回具体迭代器的实例。
迭代器的关键设计决策在于谁控制遍历——外部迭代器(如Java的Iterator)由客户端控制,显式调用next();内部迭代器(如函数式语言的forEach)由迭代器控制,客户端提供操作函数。现代编程中,两种形态常结合使用。
三、深度案例:企业级大数据查询引擎
以下展示一个真实场景下的迭代器模式应用——金融数据平台的分布式查询引擎,支持跨数据源、跨分片、跨存储介质的统一遍历抽象。
3.1 问题域分析:异构数据的统一遍历
java
// 反模式:针对不同数据源的重复遍历逻辑
public class NaiveDataQueryService {
// 遍历MySQL结果
public List<Transaction> queryFromMySql(QueryCriteria criteria) {
List<Transaction> results = new ArrayList<>();
try (Connection conn = dataSource.getConnection();
PreparedStatement stmt = conn.prepareStatement(sql);
ResultSet rs = stmt.executeQuery()) {
while (rs.next()) {
Transaction t = new Transaction();
t.setId(rs.getString("id"));
t.setAmount(rs.getBigDecimal("amount"));
// ... 更多字段映射
results.add(t);
}
}
return results;
}
// 遍历MongoDB结果
public List<Transaction> queryFromMongo(QueryCriteria criteria) {
List<Transaction> results = new ArrayList<>();
MongoCursor<Document> cursor = collection.find(filter).iterator();
while (cursor.hasNext()) {
Document doc = cursor.next();
Transaction t = new Transaction();
t.setId(doc.getString("_id"));
t.setAmount(new BigDecimal(doc.getString("amount")));
// ... 不同映射逻辑
results.add(t);
}
return results;
}
// 遍历Kafka流
public List<Transaction> queryFromKafka(String topic, QueryCriteria criteria) {
List<Transaction> results = new ArrayList<>();
KafkaConsumer<String, String> consumer = createConsumer();
consumer.subscribe(Collections.singletonList(topic));
while (results.size() < criteria.getLimit()) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
Transaction t = deserialize(record.value());
results.add(t);
}
}
return results;
}
// 每种数据源不同遍历方式,无法统一处理、无法惰性求值、无法流式处理
}
上述代码存在严重问题:遍历逻辑与数据源强耦合、全量加载内存、无法统一处理中间操作(过滤、映射、聚合)。迭代器模式通过抽象遍历接口,彻底化解这些困境。
3.2 抽象层:统一迭代器接口
java
/**
* 抽象迭代器:大数据查询迭代器
* 支持惰性加载、分页获取、资源管理、性能监控
*/
public interface QueryIterator<T> extends AutoCloseable {
/**
* 是否还有更多元素
*/
boolean hasNext();
/**
* 获取下一个元素
*/
T next();
/**
* 获取下一个元素(带默认值)
*/
default T nextOrDefault(T defaultValue) {
return hasNext() ? next() : defaultValue;
}
/**
* 跳过指定数量元素
*/
default QueryIterator<T> skip(int count) {
for (int i = 0; i < count && hasNext(); i++) {
next();
}
return this;
}
/**
* 限制返回数量
*/
default QueryIterator<T> limit(int maxCount) {
return new LimitIterator<>(this, maxCount);
}
/**
* 映射转换
*/
default <R> QueryIterator<R> map(Function<T, R> mapper) {
return new MappingIterator<>(this, mapper);
}
/**
* 过滤
*/
default QueryIterator<T> filter(Predicate<T> predicate) {
return new FilteringIterator<>(this, predicate);
}
/**
* 获取预估总数(可能不准确,用于进度展示)
*/
long estimatedTotal();
/**
* 获取当前进度
*/
default double progress() {
return -1; // 未知
}
/**
* 获取遍历统计
*/
IteratorMetrics getMetrics();
/**
* 转换为Stream(Java 8+集成)
*/
default Stream<T> toStream() {
Spliterator<T> spliterator = new IteratorSpliterator<>(this, estimatedTotal());
return StreamSupport.stream(spliterator, false)
.onClose(this::close);
}
/**
* 批量获取(优化网络往返)
*/
default List<T> nextBatch(int batchSize) {
List<T> batch = new ArrayList<>(Math.min(batchSize, 1000));
for (int i = 0; i < batchSize && hasNext(); i++) {
batch.add(next());
}
return batch;
}
@Override
void close(); // 资源释放
}
/**
* 迭代器装饰基类:支持中间操作的组合
*/
abstract class IteratorDecorator<T> implements QueryIterator<T> {
protected final QueryIterator<T> delegate;
IteratorDecorator(QueryIterator<T> delegate) {
this.delegate = delegate;
}
@Override
public long estimatedTotal() {
return delegate.estimatedTotal();
}
@Override
public double progress() {
return delegate.progress();
}
@Override
public IteratorMetrics getMetrics() {
return delegate.getMetrics();
}
@Override
public void close() {
delegate.close();
}
}
/**
* 具体装饰:限制迭代器
*/
class LimitIterator<T> extends IteratorDecorator<T> {
private final int limit;
private int count;
LimitIterator(QueryIterator<T> delegate, int limit) {
super(delegate);
this.limit = limit;
}
@Override
public boolean hasNext() {
return count < limit && delegate.hasNext();
}
@Override
public T next() {
if (!hasNext()) {
throw new NoSuchElementException();
}
count++;
return delegate.next();
}
@Override
public long estimatedTotal() {
return Math.min(limit, delegate.estimatedTotal());
}
}
/**
* 具体装饰:映射迭代器
*/
class MappingIterator<T, R> implements QueryIterator<R> {
private final QueryIterator<T> delegate;
private final Function<T, R> mapper;
MappingIterator(QueryIterator<T> delegate, Function<T, R> mapper) {
this.delegate = delegate;
this.mapper = mapper;
}
@Override
public boolean hasNext() {
return delegate.hasNext();
}
@Override
public R next() {
return mapper.apply(delegate.next());
}
@Override
public long estimatedTotal() {
return delegate.estimatedTotal();
}
@Override
public IteratorMetrics getMetrics() {
return delegate.getMetrics();
}
@Override
public void close() {
delegate.close();
}
}
/**
* 具体装饰:过滤迭代器
*/
class FilteringIterator<T> extends IteratorDecorator<T> {
private final Predicate<T> predicate;
private T nextElement;
private boolean hasNextElement;
FilteringIterator(QueryIterator<T> delegate, Predicate<T> predicate) {
super(delegate);
this.predicate = predicate;
advance();
}
private void advance() {
while (delegate.hasNext()) {
T candidate = delegate.next();
if (predicate.test(candidate)) {
nextElement = candidate;
hasNextElement = true;
return;
}
}
hasNextElement = false;
}
@Override
public boolean hasNext() {
return hasNextElement;
}
@Override
public T next() {
if (!hasNextElement) {
throw new NoSuchElementException();
}
T result = nextElement;
advance();
return result;
}
@Override
public long estimatedTotal() {
// 过滤后数量未知,返回保守估计
return delegate.estimatedTotal() / 2;
}
}
3.3 具体迭代器:异构数据源实现
java
/**
* 具体迭代器:MySQL分页查询迭代器
* 自动管理JDBC资源,后台预加载
*/
public class MySqlPageIterator<T> implements QueryIterator<T> {
private final Connection connection;
private final PreparedStatement statement;
private final ResultSet resultSet;
private final RowMapper<T> rowMapper;
private final int pageSize;
// 分页状态
private final String baseSql;
private final List<Object> parameters;
private long currentOffset;
private int currentPageRow;
private List<T> currentPage;
private boolean hasMorePages;
// 性能监控
private final IteratorMetrics metrics = new IteratorMetrics();
private final StopWatch stopWatch = new StopWatch();
MySqlPageIterator(DataSource dataSource, String sql, List<Object> params,
RowMapper<T> rowMapper, int pageSize) throws SQLException {
this.connection = dataSource.getConnection();
this.baseSql = sql;
this.parameters = params;
this.rowMapper = rowMapper;
this.pageSize = pageSize;
this.currentOffset = 0;
this.currentPageRow = 0;
// 初始加载第一页
loadNextPage();
}
@Override
public boolean hasNext() {
if (currentPageRow < currentPage.size()) {
return true;
}
if (hasMorePages) {
loadNextPage();
return currentPageRow < currentPage.size();
}
return false;
}
@Override
public T next() {
if (!hasNext()) {
throw new NoSuchElementException();
}
metrics.incrementReturned();
return currentPage.get(currentPageRow++);
}
private void loadNextPage() {
stopWatch.start();
String pageSql = baseSql + " LIMIT ? OFFSET ?";
List<Object> pageParams = new ArrayList<>(parameters);
pageParams.add(pageSize);
pageParams.add(currentOffset);
try {
if (statement != null) statement.close();
PreparedStatement stmt = connection.prepareStatement(pageSql);
for (int i = 0; i < pageParams.size(); i++) {
stmt.setObject(i + 1, pageParams.get(i));
}
this.statement = stmt;
this.resultSet = stmt.executeQuery();
currentPage = new ArrayList<>();
while (resultSet.next()) {
currentPage.add(rowMapper.mapRow(resultSet, currentPage.size()));
}
currentPageRow = 0;
currentOffset += currentPage.size();
hasMorePages = currentPage.size() == pageSize;
metrics.addPageLoadTime(stopWatch.getLastTaskTimeMillis());
} catch (SQLException e) {
throw new DataAccessException("分页加载失败", e);
}
}
@Override
public long estimatedTotal() {
// 执行COUNT查询获取预估
try (PreparedStatement countStmt = connection.prepareStatement(
"SELECT COUNT(*) FROM (" + baseSql + ") t")) {
for (int i = 0; i < parameters.size(); i++) {
countStmt.setObject(i + 1, parameters.get(i));
}
ResultSet rs = countStmt.executeQuery();
rs.next();
return rs.getLong(1);
} catch (SQLException e) {
return -1;
}
}
@Override
public double progress() {
long total = estimatedTotal();
if (total <= 0) return -1;
return (double) (currentOffset - currentPage.size() + currentPageRow) / total;
}
@Override
public IteratorMetrics getMetrics() {
return metrics;
}
@Override
public void close() {
try {
if (resultSet != null) resultSet.close();
if (statement != null) statement.close();
if (connection != null) connection.close();
} catch (SQLException e) {
throw new DataAccessException("关闭资源失败", e);
}
}
}
/**
* 具体迭代器:MongoDB游标迭代器
* 利用MongoDB原生游标的惰性加载
*/
public class MongoCursorIterator<T> implements QueryIterator<T> {
private final MongoCursor<Document> cursor;
private final DocumentMapper<T> mapper;
private final MongoCollection<Document> collection;
private final Bson filter;
// 预读缓冲
private final ArrayDeque<T> buffer = new ArrayDeque<>();
private static final int PRE_FETCH_SIZE = 100;
MongoCursorIterator(MongoCollection<Document> collection, Bson filter,
DocumentMapper<T> mapper) {
this.collection = collection;
this.filter = filter;
this.mapper = mapper;
this.cursor = collection.find(filter).batchSize(PRE_FETCH_SIZE).iterator();
preFetch();
}
private void preFetch() {
while (buffer.size() < PRE_FETCH_SIZE && cursor.hasNext()) {
buffer.add(mapper.map(cursor.next()));
}
}
@Override
public boolean hasNext() {
return !buffer.isEmpty();
}
@Override
public T next() {
if (buffer.isEmpty()) {
throw new NoSuchElementException();
}
T result = buffer.poll();
if (buffer.size() < PRE_FETCH_SIZE / 2) {
preFetch(); // 后台预加载
}
return result;
}
@Override
public long estimatedTotal() {
return collection.countDocuments(filter);
}
@Override
public void close() {
cursor.close();
}
}
/**
* 具体迭代器:Kafka流式迭代器
* 支持消费者组、偏移量管理、再平衡
*/
public class KafkaStreamIterator<T> implements QueryIterator<T> {
private final KafkaConsumer<String, String> consumer;
private final String topic;
private final Duration pollTimeout;
private final Deserializer<T> deserializer;
private ConsumerRecords<String, String> currentRecords;
private Iterator<ConsumerRecord<String, String>> currentIterator;
private boolean noMoreRecords = false;
// 偏移量管理
private final Map<TopicPartition, Long> committedOffsets = new HashMap<>();
KafkaStreamIterator(KafkaConsumer<String, String> consumer, String topic,
Duration pollTimeout, Deserializer<T> deserializer) {
this.consumer = consumer;
this.topic = topic;
this.pollTimeout = pollTimeout;
this.deserializer = deserializer;
consumer.subscribe(Collections.singletonList(topic),
new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
commitOffsets();
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
// 恢复偏移量
partitions.forEach(tp -> {
Long offset = committedOffsets.get(tp);
if (offset != null) {
consumer.seek(tp, offset);
}
});
}
});
pollNext();
}
private void pollNext() {
if (noMoreRecords) return;
currentRecords = consumer.poll(pollTimeout);
currentIterator = currentRecords.iterator();
if (currentRecords.isEmpty()) {
// 检查是否到达末尾
Set<TopicPartition> assignment = consumer.assignment();
Map<TopicPartition, Long> endOffsets = consumer.endOffsets(assignment);
for (TopicPartition tp : assignment) {
long position = consumer.position(tp);
if (position < endOffsets.get(tp)) {
return; // 还有更多数据
}
}
noMoreRecords = true;
}
}
@Override
public boolean hasNext() {
if (currentIterator != null && currentIterator.hasNext()) {
return true;
}
if (noMoreRecords) return false;
pollNext();
return currentIterator.hasNext();
}
@Override
public T next() {
if (!hasNext()) {
throw new NoSuchElementException();
}
ConsumerRecord<String, String> record = currentIterator.next();
// 记录偏移量(稍后提交)
committedOffsets.put(
new TopicPartition(record.topic(), record.partition()),
record.offset() + 1);
return deserializer.deserialize(record.value());
}
@Override
public void close() {
commitOffsets();
consumer.close();
}
private void commitOffsets() {
if (!committedOffsets.isEmpty()) {
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
committedOffsets.forEach((tp, offset) ->
offsets.put(tp, new OffsetAndMetadata(offset)));
consumer.commitSync(offsets);
committedOffsets.clear();
}
}
}
/**
* 具体迭代器:分布式聚合结果迭代器
* 合并多个数据分片的查询结果,支持排序归并
*/
public class ShardedMergeIterator<T> implements QueryIterator<T> {
private final PriorityQueue<ShardCursor<T>> heap;
private final Comparator<T> comparator;
private final List<QueryIterator<T>> shards;
ShardedMergeIterator(List<QueryIterator<T>> shards, Comparator<T> comparator) {
this.shards = shards;
this.comparator = comparator;
this.heap = new PriorityQueue<>(Comparator.comparing(
ShardCursor::getCurrent, comparator));
// 初始化堆:每个分片取第一个元素
for (QueryIterator<T> shard : shards) {
if (shard.hasNext()) {
heap.offer(new ShardCursor<>(shard, shard.next()));
}
}
}
@Override
public boolean hasNext() {
return !heap.isEmpty();
}
@Override
public T next() {
if (heap.isEmpty()) {
throw new NoSuchElementException();
}
// 取出最小元素
ShardCursor<T> min = heap.poll();
T result = min.getCurrent();
// 从同一分片取下一个元素入堆
if (min.getIterator().hasNext()) {
min.advance();
heap.offer(min);
}
return result;
}
@Override
public long estimatedTotal() {
return shards.stream().mapToLong(QueryIterator::estimatedTotal).sum();
}
@Override
public void close() {
shards.forEach(QueryIterator::close);
}
private static class ShardCursor<T> {
private final QueryIterator<T> iterator;
private T current;
ShardCursor(QueryIterator<T> iterator, T current) {
this.iterator = iterator;
this.current = current;
}
void advance() {
this.current = iterator.next();
}
T getCurrent() { return current; }
QueryIterator<T> getIterator() { return iterator; }
}
}
3.4 聚合对象:查询结果集
java
/**
* 抽象聚合:查询结果集
*/
public interface QueryResultSet<T> extends Iterable<T> {
QueryIterator<T> iterator();
/**
* 获取结果集元数据
*/
ResultSetMetadata getMetadata();
/**
* 是否支持随机访问
*/
boolean supportsRandomAccess();
/**
* 获取指定位置元素(如支持)
*/
default T get(long index) {
throw new UnsupportedOperationException("不支持随机访问");
}
/**
* 转换为List(慎用,可能内存溢出)
*/
default List<T> toList() {
List<T> list = new ArrayList<>();
try (QueryIterator<T> it = iterator()) {
while (it.hasNext()) {
list.add(it.next());
}
}
return list;
}
}
/**
* 具体聚合:SQL查询结果集
*/
public class SqlQueryResultSet<T> implements QueryResultSet<T> {
private final DataSource dataSource;
private final String sql;
private final List<Object> parameters;
private final RowMapper<T> rowMapper;
private final int pageSize;
public SqlQueryResultSet(DataSource dataSource, String sql,
List<Object> parameters,
RowMapper<T> rowMapper,
int pageSize) {
this.dataSource = dataSource;
this.sql = sql;
this.parameters = parameters;
this.rowMapper = rowMapper;
this.pageSize = pageSize;
}
@Override
public QueryIterator<T> iterator() {
try {
return new MySqlPageIterator<>(dataSource, sql, parameters,
rowMapper, pageSize);
} catch (SQLException e) {
throw new DataAccessException("创建迭代器失败", e);
}
}
@Override
public ResultSetMetadata getMetadata() {
// 解析SQL获取列信息
return SqlMetadataParser.parse(sql);
}
@Override
public boolean supportsRandomAccess() {
return false; // 流式结果不支持
}
}
/**
* 具体聚合:跨分片聚合结果集
*/
public class ShardedQueryResultSet<T> implements QueryResultSet<T> {
private final List<QueryResultSet<T>> shards;
private final Comparator<T> sortComparator;
public ShardedQueryResultSet(List<QueryResultSet<T>> shards,
Comparator<T> sortComparator) {
this.shards = shards;
this.sortComparator = sortComparator;
}
@Override
public QueryIterator<T> iterator() {
List<QueryIterator<T>> shardIterators = shards.stream()
.map(QueryResultSet::iterator)
.collect(Collectors.toList());
return new ShardedMergeIterator<>(shardIterators, sortComparator);
}
@Override
public ResultSetMetadata getMetadata() {
// 取第一个分片的元数据(假设同构)
return shards.get(0).getMetadata();
}
@Override
public boolean supportsRandomAccess() {
return false;
}
}
3.5 客户端使用:统一查询接口
java
/**
* 查询服务:客户端面向统一接口,无需关心底层数据源
*/
@Service
public class UnifiedQueryService {
private final QueryRouter queryRouter;
private final ResultCache resultCache;
/**
* 执行查询,返回惰性迭代结果
*/
public <T> QueryResultSet<T> execute(QueryRequest<T> request) {
// 路由到合适的执行器
QueryExecutor<T> executor = queryRouter.resolve(request);
// 尝试缓存
if (request.isCacheable()) {
QueryResultSet<T> cached = resultCache.get(request.getCacheKey());
if (cached != null) return cached;
}
// 执行查询
QueryResultSet<T> result = executor.execute(request);
// 缓存结果
if (request.isCacheable()) {
resultCache.put(request.getCacheKey(), result, request.getCacheTtl());
}
return result;
}
/**
* 流式处理:大数据量导出
*/
public void exportToStream(QueryRequest<Transaction> request,
OutputStream outputStream) {
QueryResultSet<Transaction> resultSet = execute(request);
try (QueryIterator<Transaction> iterator = resultSet.iterator();
JsonGenerator jsonGen = new JsonFactory().createGenerator(outputStream)) {
jsonGen.writeStartArray();
while (iterator.hasNext()) {
Transaction tx = iterator.next();
jsonGen.writeObject(tx);
// 定期刷新,避免内存缓冲
if (iterator.getMetrics().getReturnedCount() % 1000 == 0) {
jsonGen.flush();
}
}
jsonGen.writeEndArray();
} catch (IOException e) {
throw new ExportException("流式导出失败", e);
}
}
/**
* 聚合计算:利用迭代器避免全量加载
*/
public BigDecimal aggregateAmount(QueryRequest<Transaction> request) {
QueryResultSet<Transaction> resultSet = execute(request);
try (QueryIterator<Transaction> iterator = resultSet.iterator()) {
return iterator.toStream()
.map(Transaction::getAmount)
.reduce(BigDecimal.ZERO, BigDecimal::add);
}
}
/**
* 分页展示:利用迭代器的skip/limit
*/
public Page<Transaction> queryPage(QueryRequest<Transaction> request,
int pageNumber, int pageSize) {
QueryResultSet<Transaction> resultSet = execute(request);
try (QueryIterator<Transaction> iterator = resultSet.iterator()
.skip((long) pageNumber * pageSize)
.limit(pageSize)) {
List<Transaction> content = new ArrayList<>();
while (iterator.hasNext()) {
content.add(iterator.next());
}
return Page.of(content, pageNumber, pageSize,
resultSet.getMetadata().getEstimatedTotal());
}
}
}
四、迭代器模式的高级主题
4.1 响应式迭代器:Reactive Streams
java
/**
* 响应式迭代器:背压感知的异步遍历
*/
public class ReactiveQueryIterator<T> implements QueryIterator<T> {
private final Flux<T> flux;
private final Iterator<T> blockingIterator;
ReactiveQueryIterator(Flux<T> flux) {
this.flux = flux;
// 阻塞迭代适配(仅用于兼容)
this.blockingIterator = flux.toIterable(1).iterator();
}
@Override
public boolean hasNext() {
return blockingIterator.hasNext();
}
@Override
public T next() {
return blockingIterator.next();
}
/**
* 获取响应式流(推荐方式)
*/
public Flux<T> toFlux() {
return flux;
}
}
// 使用示例
public Mono<Long> reactiveAggregate(QueryRequest<Transaction> request) {
return queryService.executeReactive(request)
.toFlux()
.map(Transaction::getAmount)
.reduce(BigDecimal.ZERO, BigDecimal::add)
.map(BigDecimal::longValue);
}
4.2 并行迭代器:多线程分片处理
java
/**
* 并行迭代器:利用ForkJoinPool并行处理
*/
public class ParallelQueryIterator<T> implements QueryIterator<T> {
private final ForkJoinPool executor;
private final BlockingQueue<T> resultQueue;
private final List<Future<?>> futures;
ParallelQueryIterator(QueryIterator<T> source, int parallelism,
Function<T, T> processor) {
this.executor = new ForkJoinPool(parallelism);
this.resultQueue = new LinkedBlockingQueue<>();
// 提交并行处理任务
this.futures = IntStream.range(0, parallelism)
.mapToObj(i -> executor.submit(() -> {
while (source.hasNext()) {
T item = source.next();
T processed = processor.apply(item);
resultQueue.put(processed);
}
}))
.collect(Collectors.toList());
}
@Override
public boolean hasNext() {
return !resultQueue.isEmpty() || futures.stream().anyMatch(f -> !f.isDone());
}
@Override
public T next() {
try {
return resultQueue.take();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IterationInterruptedException(e);
}
}
@Override
public void close() {
futures.forEach(f -> f.cancel(true));
executor.shutdown();
}
}
五、迭代器模式与相关模式的辨析
迭代器 vs 访问者:迭代器遍历元素,访问者对元素执行操作。二者可结合:迭代器遍历,访问者处理。
迭代器 vs 组合模式:组合模式构建树形结构,迭代器提供遍历方式。树的深度优先、广度优先遍历可通过不同迭代器实现。
迭代器 vs 生成器:生成器(如Python的yield)是迭代器的语法糖,更简洁地实现惰性求值。
六、设计陷阱与规避策略
陷阱一:并发修改异常
遍历过程中集合被修改。解决方案:使用并发集合、CopyOnWrite策略、或显式的版本控制。
陷阱二:资源泄漏
迭代器未关闭导致连接泄漏。解决方案:实现AutoCloseable,强制try-with-resources使用。
陷阱三:全量加载伪装成惰性迭代
toList()等方法的滥用。解决方案:API设计区分惰性操作与终止操作,文档化内存影响。
七、结语
迭代器模式是遍历抽象的理论基石,它将数据结构的访问方式与数据结构本身解耦,使客户端能够以统一、安全、高效的方式遍历各种异构数据。在大数据查询、流式处理、分布式计算等场景中,迭代器模式是不可或缺的基础设施。理解其惰性求值、装饰组合、资源管理等高级机制,掌握与响应式编程、并行计算的融合,警惕并发修改、资源泄漏、全量加载等工程陷阱,是构建高性能数据系统的能力核心。迭代器模式的精髓在于尊重数据的海量本质,以受控的、增量的、可组合的方式访问数据,在有限资源与无限数据之间寻找最优的平衡点。