DATAGEODE-151 - Polish.

Resolves gh-6.
This commit is contained in:
John Blum
2018-12-18 17:53:48 -08:00
parent d26b6c3466
commit d4eab653cf
2 changed files with 302 additions and 192 deletions

View File

@@ -12,164 +12,180 @@
*/ */
package org.springframework.data.gemfire.function; package org.springframework.data.gemfire.function;
import java.lang.reflect.Array;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collection;
import java.util.Iterator;
import java.util.List;
import org.apache.geode.cache.execute.Function;
import org.apache.geode.cache.execute.ResultSender; import org.apache.geode.cache.execute.ResultSender;
import org.springframework.util.Assert; import org.springframework.util.Assert;
import org.springframework.util.ObjectUtils; import org.springframework.util.ObjectUtils;
import java.lang.reflect.Array;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Iterator;
import java.util.List;
/** /**
* Sends collection results using a {@link ResultSender} in chunks determined by batchSize * Sends {@link Collection} {@link Function} results using a {@link ResultSender} in chunks
* determined by {@code batchSize}.
* *
* @author David Turanski * @author David Turanski
* @author Udo Kohlmeyer * @author Udo Kohlmeyer
* @author John Blum
* @since 1.3.0 * @since 1.3.0
*/ */
class BatchingResultSender { class BatchingResultSender {
private final int batchSize;
private ResultSender<Object> resultSender;
public BatchingResultSender(int batchSize, ResultSender<Object> resultSender) { private final int batchSize;
Assert.notNull(resultSender, "resultSender cannot be null");
Assert.isTrue(batchSize >= 0, "batchSize must be >= 0");
this.batchSize = batchSize;
this.resultSender = resultSender;
}
private ResultSender<Object> resultSender;
public void sendResults(Iterable<?> result) { /**
//Handle both batchSize == 0 or empty Iterable * Constructs a new instance of {@link BatchingResultSender} initialized with the given {@link Integer batch size}
if (batchSize == 0 || !result.iterator().hasNext()) { * and {@link ResultSender} object used to delegate all send operations.
resultSender.lastResult(result); *
return; * @param batchSize {@link Integer} specifying the configured batch size.
} * @param resultSender {@link ResultSender} used to delegate all send operations.
* @throws IllegalArgumentException if {@link ResultSender} is {@literal null}
* or {@code batchSize} is less than {@literal 0}.
* @see org.apache.geode.cache.execute.ResultSender
*/
public BatchingResultSender(int batchSize, ResultSender<Object> resultSender) {
List<Object> chunk = new ArrayList<Object>(batchSize); Assert.notNull(resultSender, "ResultSender must not be null");
Assert.isTrue(batchSize >= 0, "batchSize must be greater than equal to 0");
for (Iterator<?> it = result.iterator(); it.hasNext(); ) { this.batchSize = batchSize;
if (chunk.size() < batchSize) { this.resultSender = resultSender;
chunk.add(it.next()); }
}
if (chunk.size() == batchSize || !it.hasNext()) { /**
if (it.hasNext()) { * Returns the configured {@link Integer batchSize} of this batching {@link ResultSender}.
resultSender.sendResult(chunk); *
} else { * @return an {@link Integer} value specifying the configured {@link Integer batchSize}
resultSender.lastResult(chunk); * of this batching {@link ResultSender}.
} */
chunk.clear(); public int getBatchSize() {
} return this.batchSize;
} }
}
/**
* Returns a reference to the configured {@link ResultSender} used to send {@link Function} results.
*
* @return a reference to the configured {@link ResultSender} used to send {@link Function} results.
* @see org.apache.geode.cache.execute.ResultSender
*/
public ResultSender<Object> getResultSender() {
return this.resultSender;
}
public void sendArrayResults(Object result) { protected boolean isBatchingDisabled() {
Assert.isTrue(ObjectUtils.isArray(result)); return !isBatchingEnabled();
}
if (batchSize == 0) { protected boolean isBatchingEnabled() {
resultSender.lastResult(result); return getBatchSize() > 0;
return; }
}
int length = Array.getLength(result); protected boolean doNotSendChunks(boolean resultSetIsEmpty) {
return resultSetIsEmpty || isBatchingDisabled();
}
if (length == 0) { public void sendResults(Iterable<?> result) {
resultSender.lastResult(result);
return;
}
for (int from = 0; from < length; from += batchSize) { ResultSender<Object> resultSender = getResultSender();
int to = Math.min(length, from + batchSize);
Object chunk = copyOfRange(result, from, to);
if (to == length) { if (doNotSendChunks(!result.iterator().hasNext())) {
resultSender.lastResult(chunk); resultSender.lastResult(result);
} else { }
resultSender.sendResult(chunk); else {
}
}
}
int batchSize = getBatchSize();
/** List<Object> chunk = new ArrayList<>(batchSize);
* @param result
* @param from
* @param to
* @return
*/
private Object copyOfRange(Object result, int from, int to) {
Class<?> arrayClass = result.getClass();
int size = to - from;
if (int[].class.isAssignableFrom(arrayClass)) { for (Iterator<?> it = result.iterator(); it.hasNext(); ) {
int[] array = new int[size];
for (int i = 0; i < size; ++i) {
array[i] = Array.getInt(result, from + i);
}
return array;
}
if (float[].class.isAssignableFrom(arrayClass)) { if (chunk.size() < batchSize) {
float[] array = new float[size]; chunk.add(it.next());
for (int i = 0; i < size; ++i) { }
array[i] = Array.getFloat(result, from + i);
}
return array;
}
if (double[].class.isAssignableFrom(arrayClass)) { if (chunk.size() == batchSize || !it.hasNext()) {
double[] array = new double[size];
for (int i = 0; i < size; ++i) {
array[i] = Array.getDouble(result, from + i);
}
return array;
}
if (boolean[].class.isAssignableFrom(arrayClass)) { if (it.hasNext()) {
boolean[] array = new boolean[size]; resultSender.sendResult(chunk);
for (int i = 0; i < size; ++i) { }
array[i] = Array.getBoolean(result, from + i); else {
} resultSender.lastResult(chunk);
return array; }
}
if (byte[].class.isAssignableFrom(arrayClass)) { chunk.clear();
byte[] array = new byte[size]; }
for (int i = 0; i < size; ++i) { }
array[i] = Array.getByte(result, from + i); }
} }
return array;
}
if (short[].class.isAssignableFrom(arrayClass)) { public void sendArrayResults(Object result) {
short[] array = new short[size];
for (int i = 0; i < size; ++i) {
array[i] = Array.getShort(result, from + i);
}
return array;
}
if (long[].class.isAssignableFrom(arrayClass)) { Assert.isTrue(ObjectUtils.isArray(result),
long[] array = new long[size]; () -> String.format("Object must be an array; was [%s]", ObjectUtils.nullSafeClassName(result)));
for (int i = 0; i < size; ++i) {
array[i] = Array.getLong(result, from + i);
}
return array;
}
if (char[].class.isAssignableFrom(arrayClass)) { int arrayLength = Array.getLength(result);
char[] array = new char[size];
for (int i = 0; i < size; ++i) {
array[i] = Array.getChar(result, from + i);
}
return array;
}
return Arrays.copyOfRange((Object[]) result, from, to); ResultSender<Object> resultSender = getResultSender();
} if (doNotSendChunks(arrayLength == 0)) {
resultSender.lastResult(result);
}
else {
int batchSize = getBatchSize();
for (int from = 0; from < arrayLength; from += batchSize) {
int to = Math.min(arrayLength, from + batchSize);
Object chunk = copyOfRange(result, from, to);
if (to == arrayLength) {
resultSender.lastResult(chunk);
}
else {
resultSender.sendResult(chunk);
}
}
}
}
private Object copyOfRange(Object result, int from, int to) {
Class<?> resultType = result.getClass();
if (boolean[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((boolean[]) result, from, to);
}
else if (byte[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((byte[]) result, from, to);
}
else if (short[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((short[]) result, from, to);
}
else if (int[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((int[]) result, from, to);
}
else if (long[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((long[]) result, from, to);
}
else if (float[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((float[]) result, from, to);
}
else if (double[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((double[]) result, from, to);
}
else if (char[].class.isAssignableFrom(resultType)) {
return Arrays.copyOfRange((char[]) result, from, to);
}
else {
return Arrays.copyOfRange((Object[]) result, from, to);
}
}
} }

View File

@@ -12,76 +12,179 @@
*/ */
package org.springframework.data.gemfire.function; package org.springframework.data.gemfire.function;
import static org.junit.Assert.assertEquals; import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.mock;
import static org.junit.Assert.fail;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Collection; import java.util.Collection;
import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Optional;
import java.util.stream.IntStream;
import org.apache.geode.cache.execute.ResultSender; import org.apache.geode.cache.execute.ResultSender;
import org.assertj.core.api.Assertions;
import org.junit.Test; import org.junit.Test;
/** /**
* Unit tests for {@link BatchingResultSender}.
*
* @author David Turanski * @author David Turanski
* @author Udo Kohlmeyer * @author Udo Kohlmeyer
* * @author John Blum
* @see org.junit.Test
* @see org.apache.geode.cache.execute.ResultSender
* @see org.springframework.data.gemfire.function.BatchingResultSender
* @since 1.3.0
*/ */
public class BatchingResultSenderTest { public class BatchingResultSenderTest {
@Test @Test
public void testArrayChunking() { @SuppressWarnings("unchecked")
public void constructBatchingResultSender() {
ResultSender<Object> mockResultSender = mock(ResultSender.class);
BatchingResultSender batchResultSender = new BatchingResultSender(20, mockResultSender);
assertThat(batchResultSender).isNotNull();
assertThat(batchResultSender.getBatchSize()).isEqualTo(20);
assertThat(batchResultSender.getResultSender()).isEqualTo(mockResultSender);
}
@SuppressWarnings("unchecked")
@Test(expected = IllegalArgumentException.class)
public void constructBatchingResultSenderWithBatchSizeOfMinusOne() {
try {
new BatchingResultSender(-1, mock(ResultSender.class));
}
catch (IllegalArgumentException expected) {
assertThat(expected).hasMessage("batchSize must be greater than equal to 0");
assertThat(expected).hasNoCause();
throw expected;
}
}
@Test(expected = IllegalArgumentException.class)
public void constructBatchingResultSenderWithNullResultSender() {
try {
new BatchingResultSender(0, null);
}
catch (IllegalArgumentException expected) {
assertThat(expected).hasMessage("ResultSender must not be null");
assertThat(expected).hasNoCause();
throw expected;
}
}
@SuppressWarnings("unchecked")
@Test(expected = IllegalArgumentException.class)
public void sendArrayResultsWithNonArrayThrowIllegalArgumentException() {
try {
new BatchingResultSender(20, mock(ResultSender.class)).sendArrayResults(new Object());
}
catch (IllegalArgumentException expected) {
assertThat(expected).hasMessage("Object must be an array; was [java.lang.Object]");
assertThat(expected).hasNoCause();
throw expected;
}
}
@Test
public void arrayChunkingIsCorrectWhenBatchSizeIsZero() {
testBatchingResultSender(new TestArrayResultSender(),0); // Default result set size is 100
testBatchingResultSender(new TestArrayResultSender(),0,1);
testBatchingResultSender(new TestArrayResultSender(),0,2);
}
@Test
public void arrayChunkingIsCorrectWhenBothBatchSizeAndResultSetSizeAreZero() {
testBatchingResultSender(new TestArrayResultSender(),0,0);
}
@Test
public void arrayCheckingIsCorrectWhenResultSetSizeIsZero() {
testBatchingResultSender(new TestArrayResultSender(),1,0);
testBatchingResultSender(new TestArrayResultSender(),2,0);
testBatchingResultSender(new TestArrayResultSender(),3,0);
testBatchingResultSender(new TestArrayResultSender(),50,0);
}
@Test
public void arrayChunkingIsCorrect() {
testBatchingResultSender(new TestArrayResultSender(),1); testBatchingResultSender(new TestArrayResultSender(),1);
testBatchingResultSender(new TestArrayResultSender(),0);
testBatchingResultSender(new TestArrayResultSender(),9); testBatchingResultSender(new TestArrayResultSender(),9);
testBatchingResultSender(new TestArrayResultSender(),10,99); testBatchingResultSender(new TestArrayResultSender(),10,99);
testBatchingResultSender(new TestArrayResultSender(),10,100); testBatchingResultSender(new TestArrayResultSender(),10,100);
testBatchingResultSender(new TestArrayResultSender(),10,101); testBatchingResultSender(new TestArrayResultSender(),10,101);
testBatchingResultSender(new TestArrayResultSender(),1000); testBatchingResultSender(new TestArrayResultSender(),1000);
testBatchingResultSender(new TestArrayResultSender(),2,0);
testBatchingResultSender(new TestArrayResultSender(),0,0);
} }
@Test
public void listChunkingIsCorrectWhenBatchSizeIsZero() {
testBatchingResultSender(new TestListResultSender(),0); // Default result set size is 100
testBatchingResultSender(new TestListResultSender(),0, 1);
testBatchingResultSender(new TestListResultSender(),0, 2);
}
@Test @Test
public void testListChunking() { public void listChunkingIsCorrectWhenBothBatchSizeAndResultSetSizeAreZero() {
testBatchingResultSender(new TestListResultSender(),0,0);
}
@Test
public void listChunkingIsCorrectWhenResultSetSizeIsZero() {
testBatchingResultSender(new TestListResultSender(),1,0);
testBatchingResultSender(new TestListResultSender(),2,0);
testBatchingResultSender(new TestListResultSender(),3,0);
testBatchingResultSender(new TestListResultSender(),50,0);
}
@Test
public void listChunkingIsCorrect() {
testBatchingResultSender(new TestListResultSender(),1); testBatchingResultSender(new TestListResultSender(),1);
testBatchingResultSender(new TestListResultSender(),0);
testBatchingResultSender(new TestListResultSender(),9); testBatchingResultSender(new TestListResultSender(),9);
testBatchingResultSender(new TestListResultSender(),10,99); testBatchingResultSender(new TestListResultSender(),10,99);
testBatchingResultSender(new TestListResultSender(),10,100); testBatchingResultSender(new TestListResultSender(),10,100);
testBatchingResultSender(new TestListResultSender(),10,101); testBatchingResultSender(new TestListResultSender(),10,101);
testBatchingResultSender(new TestListResultSender(),1000); testBatchingResultSender(new TestListResultSender(),1000);
testBatchingResultSender(new TestListResultSender(),3,0);
testBatchingResultSender(new TestListResultSender(),0,0);
} }
private void testBatchingResultSender(AbstractTestResultSender resultSender, int batchSize,int resultSetSize){ private void testBatchingResultSender(AbstractTestResultSender resultSender, int batchSize, int resultSetSize){
BatchingResultSender brs = new BatchingResultSender(batchSize, resultSender);
BatchingResultSender batchResultSender = new BatchingResultSender(batchSize, resultSender);
List<Integer> result = new ArrayList<>(); List<Integer> result = new ArrayList<>();
for (int i = 0; i< resultSetSize; i++) {
result.add(i); IntStream.range(0, resultSetSize).forEach(result::add);
}
//TODO: Clean this up. Ok for test code
if (resultSender instanceof TestArrayResultSender) { if (resultSender instanceof TestArrayResultSender) {
brs.sendArrayResults(result.toArray(new Integer[resultSetSize])); batchResultSender.sendArrayResults(result.toArray(new Integer[resultSetSize]));
} else { } else {
brs.sendResults(result); batchResultSender.sendResults(result);
} }
assertEquals(resultSetSize,resultSender.getResults().size()); assertThat(resultSender.isLastResultSent()).isTrue();
assertThat(resultSender.getResults()).hasSize(resultSetSize);
assertTrue(resultSender.isLastResultSent());
for(int i=0; i< resultSetSize; i++) {
assertEquals(i,resultSender.getResults().get(i));
}
IntStream.range(0, resultSetSize).forEach(index ->
assertThat(resultSender.getResults().get(index)).isEqualTo(index));
} }
private void testBatchingResultSender(AbstractTestResultSender resultSender, int batchSize){ private void testBatchingResultSender(AbstractTestResultSender resultSender, int batchSize){
@@ -89,43 +192,34 @@ public class BatchingResultSenderTest {
} }
public static abstract class AbstractTestResultSender implements ResultSender<Object> { public static abstract class AbstractTestResultSender implements ResultSender<Object> {
private List<Object> results = new ArrayList<Object>();
private boolean lastResultSent = false; private boolean lastResultSent = false;
private List<Object> results = new ArrayList<>();
public boolean isLastResultSent() { public boolean isLastResultSent() {
return lastResultSent; return this.lastResultSent;
} }
/* (non-Javadoc)
* @see org.apache.geode.cache.execute.ResultSender#lastResult(java.lang.Object)
*/
@Override @Override
public void lastResult(Object arg0) { public void lastResult(Object result) {
lastResultSent = true;
if (arg0 == null) { this.lastResultSent = true;
return;
} Optional.ofNullable(result)
addResults(arg0, results); .ifPresent(it -> addResults(it, this.results));
} }
/* (non-Javadoc)
* @see org.apache.geode.cache.execute.ResultSender#sendException(java.lang.Throwable)
*/
@Override @Override
public void sendException(Throwable arg0) { public void sendException(Throwable cause) {
fail(); Assertions.fail("Function send result operation failed", cause);
} }
/* (non-Javadoc)
* @see org.apache.geode.cache.execute.ResultSender#sendResult(java.lang.Object)
*/
@Override @Override
public void sendResult(Object arg0) { public void sendResult(Object result) {
if (arg0 == null) {
return; Optional.ofNullable(result)
} .ifPresent(it -> addResults(it, this.results));
addResults(arg0, results);
} }
protected abstract void addResults(Object item, List<Object> results); protected abstract void addResults(Object item, List<Object> results);
@@ -133,28 +227,28 @@ public class BatchingResultSenderTest {
public List<Object> getResults() { public List<Object> getResults() {
return this.results; return this.results;
} }
} }
public static class TestArrayResultSender extends AbstractTestResultSender { public static class TestArrayResultSender extends AbstractTestResultSender {
protected void addResults(Object arg0, List<Object> results) { protected void addResults(Object result, List<Object> results) {
assertTrue(arg0.getClass().isArray());
Object[] array = (Object[]) arg0; assertThat(result.getClass().isArray()).isTrue();
for (Object obj: array) {
results.add(obj); Object[] array = (Object[]) result;
}
Collections.addAll(results, array);
} }
} }
public static class TestListResultSender extends AbstractTestResultSender { public static class TestListResultSender extends AbstractTestResultSender {
protected void addResults(Object arg0, List<Object> results) {
if (arg0 == null) { protected void addResults(Object result, List<Object> results) {
return;
} assertThat(result).isInstanceOf(Collection.class);
assertTrue(arg0 instanceof Collection);
Collection<?> list = (Collection<?>) arg0; Collection<?> list = (Collection<?>) result;
results.addAll(list); results.addAll(list);
} }
} }