/
MapCombineCommand.java
206 lines (181 loc) · 6.11 KB
/
MapCombineCommand.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
/*
* JBoss, Home of Professional Open Source
* Copyright 2011 Red Hat Inc. and/or its affiliates and other
* contributors as indicated by the @author tags. All rights reserved.
* See the copyright.txt in the distribution for a full listing of
* individual contributors.
*
* This is free software; you can redistribute it and/or modify it
* under the terms of the GNU Lesser General Public License as
* published by the Free Software Foundation; either version 2.1 of
* the License, or (at your option) any later version.
*
* This software is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
* Lesser General Public License for more details.
*
* You should have received a copy of the GNU Lesser General Public
* License along with this software; if not, write to the Free
* Software Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA
* 02110-1301 USA, or see the FSF site: http://www.fsf.org.
*/
package org.infinispan.commands.read;
import java.util.Collection;
import java.util.HashSet;
import java.util.Set;
import java.util.UUID;
import org.infinispan.commands.CancellableCommand;
import org.infinispan.commands.remote.BaseRpcCommand;
import org.infinispan.context.InvocationContext;
import org.infinispan.distexec.mapreduce.MapReduceManager;
import org.infinispan.distexec.mapreduce.Mapper;
import org.infinispan.distexec.mapreduce.Reducer;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
/**
* MapCombineCommand is a container to migrate {@link Mapper} and {@link Reducer} which is a
* combiner to a remote Infinispan node where it will get executed and return the result to an
* invoking/master node.
*
* @author Vladimir Blagojevic
* @since 5.0
*/
public class MapCombineCommand<KIn, VIn, KOut, VOut> extends BaseRpcCommand implements CancellableCommand {
public static final int COMMAND_ID = 30;
private static final Log log = LogFactory.getLog(MapCombineCommand.class);
private Set<KIn> keys = new HashSet<KIn>();
private Mapper<KIn, VIn, KOut, VOut> mapper;
private Reducer<KOut, VOut> combiner;
private String taskId;
private boolean reducePhaseDistributed;
private boolean emitCompositeIntermediateKeys;
private MapReduceManager mrManager;
private UUID uuid;
public MapCombineCommand() {
super(null); // For command id uniqueness test
}
public MapCombineCommand(String cacheName) {
super(cacheName);
}
public MapCombineCommand(String taskId, Mapper<KIn, VIn, KOut, VOut> mapper,
Reducer<KOut, VOut> combiner, String cacheName, Collection<KIn> inputKeys) {
super(cacheName);
this.taskId = taskId;
if (inputKeys != null && !inputKeys.isEmpty()) {
keys.addAll(inputKeys);
}
this.mapper = mapper;
this.combiner = combiner;
this.uuid = UUID.randomUUID();
}
public void init(MapReduceManager mrManager) {
this.mrManager = mrManager;
}
/**
* Performs invocation of mapping phase and local combine phase on assigned Infinispan node
*
* @param context
* invocation context
* @return Map of intermediate key value pairs
*/
@Override
public Object perform(InvocationContext context) throws Throwable {
if (isReducePhaseDistributed())
return mrManager.mapAndCombineForDistributedReduction(this);
else
return mrManager.mapAndCombineForLocalReduction(this);
}
public boolean isEmitCompositeIntermediateKeys() {
return emitCompositeIntermediateKeys;
}
public void setEmitCompositeIntermediateKeys(boolean emitCompositeIntermediateKeys) {
this.emitCompositeIntermediateKeys = emitCompositeIntermediateKeys;
}
public boolean isReducePhaseDistributed() {
return reducePhaseDistributed;
}
public void setReducePhaseDistributed(boolean reducePhaseDistributed) {
this.reducePhaseDistributed = reducePhaseDistributed;
}
public Set<KIn> getKeys() {
return keys;
}
public Mapper<KIn, VIn, KOut, VOut> getMapper() {
return mapper;
}
public Reducer<KOut, VOut> getCombiner() {
return combiner;
}
public String getTaskId() {
return taskId;
}
@Override
public byte getCommandId() {
return COMMAND_ID;
}
@Override
public UUID getUUID() {
return uuid;
}
@Override
public Object[] getParameters() {
return new Object[] { taskId, keys, mapper, combiner, reducePhaseDistributed,
emitCompositeIntermediateKeys, uuid };
}
@SuppressWarnings("unchecked")
@Override
public void setParameters(int commandId, Object[] args) {
if (commandId != COMMAND_ID)
throw new IllegalStateException("Invalid method id");
int i = 0;
taskId = (String) args[i++];
keys = (Set<KIn>) args[i++];
mapper = (Mapper<KIn, VIn, KOut, VOut>) args[i++];
combiner = (Reducer<KOut,VOut>) args[i++];
reducePhaseDistributed = (Boolean) args[i++];
emitCompositeIntermediateKeys = (Boolean) args[i++];
uuid = (UUID) args[i++];
}
@Override
public boolean isReturnValueExpected() {
return true;
}
@Override
public boolean canBlock() {
return true;
}
@Override
public int hashCode() {
final int prime = 31;
int result = 1;
result = prime * result + ((taskId == null) ? 0 : taskId.hashCode());
return result;
}
@SuppressWarnings("rawtypes")
@Override
public boolean equals(Object obj) {
if (this == obj) {
return true;
}
if (obj == null) {
return false;
}
if (!(obj instanceof MapCombineCommand)) {
return false;
}
MapCombineCommand other = (MapCombineCommand) obj;
if (taskId == null) {
if (other.taskId != null) {
return false;
}
} else if (!taskId.equals(other.taskId)) {
return false;
}
return true;
}
@Override
public String toString() {
return "MapCombineCommand [keys=" + keys + ", taskId=" + taskId + "]";
}
}