| 122 | } |
| 123 | |
| 124 | public void fetch(RowConsumer consumer, ScrollMode scrollMode, long count) { |
| 125 | if (queryIterator.isDone()) { |
| 126 | try { |
| 127 | BatchIterator<Row> bi = queryIterator.join(); |
| 128 | triggerConsumer(consumer, new BufferingBatchIterator(bi), scrollMode, count); |
| 129 | } catch (Throwable t) { |
| 130 | consumer.accept(null, t); |
| 131 | } |
| 132 | } else { |
| 133 | queryIterator.whenComplete((bi, err) -> { |
| 134 | if (err == null) { |
| 135 | try { |
| 136 | triggerConsumer(consumer, new BufferingBatchIterator(bi), scrollMode, count); |
| 137 | } catch (Throwable t) { |
| 138 | consumer.accept(null, t); |
| 139 | } |
| 140 | } else { |
| 141 | consumer.accept(null, err); |
| 142 | } |
| 143 | }); |
| 144 | } |
| 145 | } |
| 146 | |
| 147 | private class BufferingBatchIterator extends ForwardingBatchIterator<Row> { |
| 148 | |