1
|
/**
|
2
|
* Copyright (C) 2007 EDIT
|
3
|
* European Distributed Institute of Taxonomy
|
4
|
* http://www.e-taxonomy.eu
|
5
|
*
|
6
|
* The contents of this file are subject to the Mozilla Public License Version 1.1
|
7
|
* See LICENSE.TXT at the top of this package for the full license terms.
|
8
|
*/
|
9
|
|
10
|
package eu.etaxonomy.cdm.io.common;
|
11
|
|
12
|
import java.util.ArrayList;
|
13
|
import java.util.Iterator;
|
14
|
import java.util.LinkedList;
|
15
|
import java.util.List;
|
16
|
|
17
|
import org.apache.log4j.Logger;
|
18
|
import org.springframework.transaction.TransactionStatus;
|
19
|
|
20
|
import eu.etaxonomy.cdm.api.application.ICdmRepository;
|
21
|
import eu.etaxonomy.cdm.common.monitor.IProgressMonitor;
|
22
|
import eu.etaxonomy.cdm.common.monitor.SubProgressMonitor;
|
23
|
import eu.etaxonomy.cdm.filter.TaxonNodeFilter;
|
24
|
import eu.etaxonomy.cdm.model.taxon.TaxonNode;
|
25
|
|
26
|
/**
|
27
|
* @author a.mueller
|
28
|
* @created 01.07.2017
|
29
|
*/
|
30
|
public class TaxonNodeOutStreamPartitioner<STATE extends IoStateBase> {
|
31
|
|
32
|
private static final Logger logger = Logger.getLogger(TaxonNodeOutStreamPartitioner.class);
|
33
|
|
34
|
//************************* STATIC ***************************************************/
|
35
|
|
36
|
public static <ST extends XmlExportState> TaxonNodeOutStreamPartitioner NewInstance(
|
37
|
ICdmRepository repository, IoStateBase state,
|
38
|
TaxonNodeFilter filter, Integer partitionSize,
|
39
|
IProgressMonitor parentMonitor, Integer parentTicks){
|
40
|
TaxonNodeOutStreamPartitioner<ST> taxonNodePartitioner
|
41
|
= new TaxonNodeOutStreamPartitioner(repository, state, filter, partitionSize,
|
42
|
parentMonitor, parentTicks);
|
43
|
return taxonNodePartitioner;
|
44
|
}
|
45
|
|
46
|
//*********************** VARIABLES *************************************************/
|
47
|
|
48
|
|
49
|
/**
|
50
|
* counter for the partitions
|
51
|
*/
|
52
|
private int currentPartition;
|
53
|
|
54
|
|
55
|
private TransactionStatus txStatus;
|
56
|
|
57
|
|
58
|
//******************
|
59
|
|
60
|
private final ICdmRepository repository;
|
61
|
|
62
|
private final TaxonNodeFilter filter;
|
63
|
|
64
|
private STATE state;
|
65
|
|
66
|
/**
|
67
|
* Number of records handled in the partition
|
68
|
*/
|
69
|
private final int partitionSize;
|
70
|
|
71
|
private int totalCount = -1;
|
72
|
private List<Integer> idList;
|
73
|
|
74
|
private IProgressMonitor parentMonitor;
|
75
|
private Integer parentTicks;
|
76
|
|
77
|
private SubProgressMonitor monitor;
|
78
|
|
79
|
private LinkedList<TaxonNode> fifo = new LinkedList<>();
|
80
|
|
81
|
private Iterator<Integer> idIterator;
|
82
|
|
83
|
private int currentIndex;
|
84
|
|
85
|
private static final int retrieveFactor = 1;
|
86
|
private static final int iterateFactor = 2;
|
87
|
|
88
|
|
89
|
//*********************** CONSTRUCTOR *************************************************/
|
90
|
|
91
|
private TaxonNodeOutStreamPartitioner(ICdmRepository repository, STATE state,
|
92
|
TaxonNodeFilter filter, Integer partitionSize,
|
93
|
IProgressMonitor parentMonitor, Integer parentTicks){
|
94
|
this.repository = repository;
|
95
|
this.filter = filter;
|
96
|
this.partitionSize = partitionSize;
|
97
|
this.state = state;
|
98
|
this.parentMonitor = parentMonitor;
|
99
|
this.parentTicks = parentTicks;
|
100
|
initialize();
|
101
|
|
102
|
}
|
103
|
|
104
|
//************************ METHODS ****************************************************/
|
105
|
|
106
|
public void initialize(){
|
107
|
if (totalCount < 0){
|
108
|
|
109
|
parentMonitor.subTask("Compute total number of records");
|
110
|
totalCount = ((Long)repository.getTaxonNodeService().count(filter)).intValue();
|
111
|
idList = repository.getTaxonNodeService().idList(filter);
|
112
|
int parTicks = this.parentTicks == null? totalCount : this.parentTicks;
|
113
|
|
114
|
monitor = SubProgressMonitor.NewStarted(parentMonitor, parTicks,
|
115
|
"Taxon node streamer", totalCount * (retrieveFactor + iterateFactor));
|
116
|
idIterator = idList.iterator();
|
117
|
monitor.subTask("id iterator created");
|
118
|
}
|
119
|
}
|
120
|
|
121
|
|
122
|
// public boolean hasNext() {
|
123
|
// initialize();
|
124
|
// return idIterator.hasNext()|| idList.size() > 0;
|
125
|
// }
|
126
|
|
127
|
public TaxonNode next(){
|
128
|
int currentIndexAtStart = currentIndex;
|
129
|
|
130
|
if(fifo.isEmpty()){
|
131
|
|
132
|
List<TaxonNode> list = getNextPartition();
|
133
|
fifo.addAll(list);
|
134
|
}
|
135
|
if (!fifo.isEmpty()){
|
136
|
TaxonNode result = fifo.removeFirst();
|
137
|
// worked should be called after each step is ready,
|
138
|
//this is usually after each next() call but not for the first
|
139
|
if (currentIndexAtStart > 0){
|
140
|
monitor.worked(iterateFactor);
|
141
|
}
|
142
|
return result;
|
143
|
}else{
|
144
|
commitTransaction();
|
145
|
return null;
|
146
|
}
|
147
|
}
|
148
|
|
149
|
public void close(){
|
150
|
monitor.done();
|
151
|
commitTransaction();
|
152
|
}
|
153
|
|
154
|
private List<TaxonNode> getNextPartition() {
|
155
|
List<Integer> partList = new ArrayList<>();
|
156
|
|
157
|
if (txStatus != null){
|
158
|
commitTransaction();
|
159
|
}
|
160
|
txStatus = startTransaction();
|
161
|
while (partList.size() < partitionSize && idIterator.hasNext()){
|
162
|
partList.add(idIterator.next());
|
163
|
currentIndex++;
|
164
|
}
|
165
|
List<TaxonNode> partition = new ArrayList<>();
|
166
|
if (!partList.isEmpty()){
|
167
|
monitor.subTask(String.format("Reading partition %d/%d", currentPartition + 1, (totalCount / partitionSize) +1 ));
|
168
|
List<String> propertyPaths = null;
|
169
|
partition = repository.getTaxonNodeService().loadByIds(partList, propertyPaths);
|
170
|
monitor.worked(partition.size() * 1);
|
171
|
currentPartition++;
|
172
|
monitor.subTask(String.format("Writing partition %d/%d", currentPartition, (totalCount / partitionSize) +1 ));
|
173
|
}
|
174
|
return partition;
|
175
|
}
|
176
|
|
177
|
private void commitTransaction() {
|
178
|
if (!txStatus.isCompleted()){
|
179
|
repository.commitTransaction(txStatus);
|
180
|
}
|
181
|
}
|
182
|
|
183
|
private TransactionStatus startTransaction() {
|
184
|
return repository.startTransaction();
|
185
|
}
|
186
|
|
187
|
|
188
|
|
189
|
/**
|
190
|
* @param recordsPerTransaction
|
191
|
* @param partitionedIO
|
192
|
* @param i
|
193
|
*/
|
194
|
private TransactionStatus getTransaction(int recordsPerTransaction, IPartitionedIO partitionedIO) {
|
195
|
//if (loopNeedsHandling (i, recordsPerTransaction) || txStatus == null) {
|
196
|
txStatus = partitionedIO.startTransaction();
|
197
|
if(logger.isInfoEnabled()) {
|
198
|
logger.debug("currentPartitionNumber = " + currentPartition + " - Transaction started");
|
199
|
}
|
200
|
//}
|
201
|
return txStatus;
|
202
|
}
|
203
|
|
204
|
|
205
|
|
206
|
|
207
|
|
208
|
}
|