[
https://issues.apache.org/jira/browse/DRILL-5457?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=16027261#comment-16027261
]
ASF GitHub Bot commented on DRILL-5457:
---------------------------------------
Github user paul-rogers commented on a diff in the pull request:
https://github.com/apache/drill/pull/822#discussion_r118813868
--- Diff:
exec/java-exec/src/main/java/org/apache/drill/exec/physical/impl/aggregate/HashAggTemplate.java
---
@@ -400,114 +782,411 @@ public IterOutcome getOutcome() {
@Override
public int getOutputCount() {
- // return outputCount;
return lastBatchOutputCount;
}
@Override
public void cleanup() {
- if (htable != null) {
- htable.clear();
- htable = null;
- }
+ if ( schema == null ) { return; } // not set up; nothing to clean
+ for ( int i = 0; i < numPartitions; i++) {
+ if (htables[i] != null) {
+ htables[i].clear();
+ htables[i] = null;
+ }
+ if ( batchHolders[i] != null) {
+ for (BatchHolder bh : batchHolders[i]) {
+ bh.clear();
+ }
+ batchHolders[i].clear();
+ batchHolders[i] = null;
+ }
+
+ // delete any (still active) output spill file
+ if ( outputStream[i] != null && spillFiles[i] != null) {
+ try {
+ spillSet.delete(spillFiles[i]);
+ } catch(IOException e) {
+ logger.warn("Cleanup: Failed to delete spill file
{}",spillFiles[i]);
+ }
+ }
+ }
+ // delete any spill file left in unread spilled partitions
+ while ( ! spilledPartitionsList.isEmpty() ) {
+ SpilledPartition sp = spilledPartitionsList.remove(0);
+ try {
+ spillSet.delete(sp.spillFile);
+ } catch(IOException e) {
+ logger.warn("Cleanup: Failed to delete spill file
{}",sp.spillFile);
+ }
+ }
+ spillSet.close(); // delete the spill directory(ies)
htIdxHolder = null;
materializedValueFields = null;
outStartIdxHolder = null;
outNumRecordsHolder = null;
+ }
- if (batchHolders != null) {
- for (BatchHolder bh : batchHolders) {
+ // First free the memory used by the given (spilled) partition (i.e.,
hash table plus batches)
+ // then reallocate them in pristine state to allow the partition to
continue receiving rows
+ private void reinitPartition(int part) throws SchemaChangeException,
ClassTransformationException, IOException {
--- End diff --
Method on the partition state class
> Support Spill to Disk for the Hash Aggregate Operator
> -----------------------------------------------------
>
> Key: DRILL-5457
> URL: https://issues.apache.org/jira/browse/DRILL-5457
> Project: Apache Drill
> Issue Type: Improvement
> Components: Execution - Relational Operators
> Affects Versions: 1.10.0
> Reporter: Boaz Ben-Zvi
> Assignee: Boaz Ben-Zvi
> Fix For: 1.11.0
>
>
> Support gradual spilling memory to disk as the available memory gets too
> small to allow in memory work for the Hash Aggregate Operator.
--
This message was sent by Atlassian JIRA
(v6.3.15#6346)