serialize the operations when they are handed to replication service for broadcast; still group many byte[] into one message by delayed sending

This commit is contained in:
Axel Uhl committed 2014-06-22 12:49:48 +02:00
1 parent 1b7a4e4891
commit 0a2c7fd422
4 files changed
+41 -28

No files matched your search

@@ -11,7 +11,6 @@ import java.util.UUID;
import org.junit.Before;
import org.junit.Test;
import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.operationaltransformation.CreateLeaderboardGroup;
import com.sap.sailing.server.replication.impl.ReplicaDescriptor;
import com.sap.sailing.server.replication.impl.ReplicationInstancesManager;
@@ -34,7 +33,7 @@ public class ReplicationInstancesManagerLoggingPerformanceTest {
public void testLoggingPerformance() {
final int count = 10000000;
for (int i=0; i<count; i++) {
replicationInstanceManager.log(Collections.<RacingEventServiceOperation<?>>singletonList(operation));
replicationInstanceManager.log(Collections.<Class<?>>singletonList(operation.getClass()));
}
assertEquals(count, replicationInstanceManager.getStatistics(replica).get(operation.getClass()).intValue());
assertEquals(1.0, replicationInstanceManager.getAverageNumberOfOperationsPerMessage(replica), 0.0000001);
@@ -44,7 +43,7 @@ public class ReplicationInstancesManagerLoggingPerformanceTest {
public void testLoggingAverages() {
final int count = 1000;
for (int i=0; i<count; i++) {
replicationInstanceManager.log(Arrays.asList(new RacingEventServiceOperation<?>[] { operation, operation }));
replicationInstanceManager.log(Arrays.asList(new Class<?>[] { operation.getClass(), operation.getClass() }));
}
assertEquals(2*count, replicationInstanceManager.getStatistics(replica).get(operation.getClass()).intValue());
assertEquals(2.0, replicationInstanceManager.getAverageNumberOfOperationsPerMessage(replica), 0.0000001);
@@ -3,6 +3,7 @@ package com.sap.sailing.server.replication.impl;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -80,18 +81,17 @@ public class ReplicationInstancesManager {
*
* @see #getStatistics
*/
public void log(Iterable<RacingEventServiceOperation<?>> replicatedOperations) {
public void log(List<Class<?>> classes) {
for (ReplicaDescriptor replica : getReplicaDescriptors()) {
Map<Class<? extends RacingEventServiceOperation<?>>, Integer> counts = replicationCounts.get(replica);
if (counts == null) {
counts = new HashMap<Class<? extends RacingEventServiceOperation<?>>, Integer>();
replicationCounts.put(replica, counts);
}
for (RacingEventServiceOperation<?> replicatedOperation : replicatedOperations) {
@SuppressWarnings("unchecked")
for (Class<?> replicatedOperationClass : classes) {
// safe because replicatedOperation is declared of type RacingEventserviceOperation<T>
Class<? extends RacingEventServiceOperation<?>> operationClass =
(Class<? extends RacingEventServiceOperation<?>>) replicatedOperation.getClass();
@SuppressWarnings("unchecked")
Class<? extends RacingEventServiceOperation<?>> operationClass = (Class<? extends RacingEventServiceOperation<?>>) replicatedOperationClass;
Integer count = counts.get(operationClass);
if (count == null) {
count = 0;
@@ -102,7 +102,7 @@ public class ReplicationInstancesManager {
if (totalOps == null) {
totalOps = 0l;
}
totalOps += Util.size(replicatedOperations);
totalOps += Util.size(classes);
totalNumberOfOperations.put(replica, totalOps);
Long totalMessages = totalMessageCount.get(replica);
if (totalMessages == null) {
@@ -31,6 +31,7 @@ import com.sap.sailing.server.RacingEventServiceOperation;
import com.sap.sailing.server.replication.ReplicationMasterDescriptor;
import com.sap.sailing.server.replication.ReplicationService;
import com.sap.sailing.util.BuildVersion;
import com.sap.sse.common.Util.Pair;
/**
* Can observe a {@link RacingEventService} for the operations it performs that require replication. Only observes as
@@ -92,7 +93,7 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* has passed since the last sending, managed by a {@link Timer}. Writers need to synchronize on this buffer. This includes the
* addition of an operation to the buffer as well as the atomic sending and clearing.
*/
private final List<RacingEventServiceOperation<?>> outboundBuffer;
private final List<Pair<Class<?>, byte[]>> outboundBuffer;
/**
* Used to schedule the sending of all operations in {@link #outboundBuffer} using the {@link #sendingTask}.
@@ -233,13 +234,18 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
* a {@link #timer} is created and scheduled to send in {@link #TRANSMISSION_DELAY_MILLIS} milliseconds.
*/
private void broadcastOperation(RacingEventServiceOperation<?> operation) throws IOException {
ByteArrayOutputStream bos = new ByteArrayOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(bos);
oos.writeObject(operation);
oos.close();
final byte[] bytes = bos.toByteArray();
synchronized (outboundBuffer) {
outboundBuffer.add(operation);
outboundBuffer.add(new Pair<Class<?>, byte[]>(operation.getClass(), bytes));
if (sendingTask == null) {
sendingTask = new TimerTask() {
@Override
public void run() {
final Iterable<RacingEventServiceOperation<?>> listToSend;
final Iterable<Pair<Class<?>, byte[]>> listToSend;
synchronized (outboundBuffer) {
listToSend = new ArrayList<>(outboundBuffer);
outboundBuffer.clear();
@@ -257,15 +263,21 @@ public class ReplicationServiceImpl implements ReplicationService, OperationExec
}
}
private void broadcastOperations(Iterable<RacingEventServiceOperation<?>> operationsList) throws IOException {
// serialize operation into message
ByteArrayOutputStream bos = new ByteArrayOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(bos);
oos.writeObject(operationsList);
private void broadcastOperations(Iterable<Pair<Class<?>, byte[]>> listToSend) throws IOException {
ByteArrayOutputStream buf = new ByteArrayOutputStream();
ObjectOutputStream oos = new ObjectOutputStream(buf);
List<Class<?>> classes = new ArrayList<>();
List<byte[]> byteArrays = new ArrayList<>();
for (Pair<Class<?>, byte[]> op : listToSend) {
classes.add(op.getA());
byteArrays.add(op.getB());
}
oos.writeObject(byteArrays);
oos.close();
// copy serialized operations into message
if (masterChannel != null) {
masterChannel.basicPublish(exchangeName, /* routingKey */"", /* properties */null, bos.toByteArray());
replicationInstancesManager.log(operationsList);
masterChannel.basicPublish(exchangeName, /* routingKey */"", /* properties */null, buf.toByteArray());
replicationInstancesManager.log(classes);
}
}
@@ -96,8 +96,14 @@ public class Replicator implements Runnable {
ObjectInputStream ois = racingEventServiceTracker.getRacingEventService().getBaseDomainFactory()
.createObjectInputStreamResolvingAgainstThisFactory(new ByteArrayInputStream(bytesFromMessage));
@SuppressWarnings("unchecked")
Iterable<RacingEventServiceOperation<?>> operations = (Iterable<RacingEventServiceOperation<?>>) ois.readObject();
applyOrQueue(operations);
Iterable<byte[]> byteArrays = (Iterable<byte[]>) ois.readObject();
for (byte[] serializedOperation : byteArrays) {
ObjectInputStream operationOIS = racingEventServiceTracker
.getRacingEventService().getBaseDomainFactory().createObjectInputStreamResolvingAgainstThisFactory(
new ByteArrayInputStream(serializedOperation));
RacingEventServiceOperation<?> operation = (RacingEventServiceOperation<?>) operationOIS.readObject();
applyOrQueue(operation);
}
} catch (ConsumerCancelledException cce) {
logger.info("Consumer has been shut down properly.");
break;
@@ -163,15 +169,11 @@ public class Replicator implements Runnable {
* If the replicator is currently {@link #suspended}, the <code>operation</code> is queued, otherwise immediately applied to
* the receiving replica.
*/
private synchronized void applyOrQueue(Iterable<RacingEventServiceOperation<?>> operations) {
private synchronized void applyOrQueue(RacingEventServiceOperation<?> operation) {
if (suspended) {
for (RacingEventServiceOperation<?> operation : operations) {
queue(operation);
}
queue(operation);
} else {
for (RacingEventServiceOperation<?> operation : operations) {
apply(operation);
}
apply(operation);
}
}