Class FederationUtils
java.lang.Object
org.apache.sysds.runtime.controlprogram.federated.FederationUtils
-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionstatic MatrixBlockaggAdd(Future<FederatedResponse>[] ffr) static booleanaggBooleanScalar(Future<FederatedResponse>[] tmp) static MatrixBlockaggMatrix(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map) static MatrixBlockaggMatrix(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, FederationMap map) static MatrixBlockaggMean(Future<FederatedResponse>[] ffr, FederationMap map) static MatrixBlockaggMinMax(Future<FederatedResponse>[] ffr, boolean isMin, boolean isScalar, Optional<FTypes.FType> fedType) static MatrixBlockaggMinMaxIndex(Future<FederatedResponse>[] ffr, boolean isMin, FederationMap map) static MatrixBlockaggProd(Future<FederatedResponse>[] ffr, FederationMap fedMap, AggregateUnaryOperator aop) static MatrixBlockaggregateResponses(List<org.apache.commons.lang3.tuple.Pair<FederatedRange, Future<FederatedResponse>>> readResponses) Aggregate partially aggregated data from federated workers by adding values with the same index in different federated locations.static ScalarObjectaggScalar(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr) static ScalarObjectaggScalar(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map) static ScalarObjectaggScalar(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, FederationMap map) static MatrixBlockaggVar(Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map, boolean isRowAggregate, boolean isScalar) static MatrixBlockbind(Future<FederatedResponse>[] ffr, boolean cbind) static MatrixBlockbindResponses(List<org.apache.commons.lang3.tuple.Pair<FederatedRange, Future<FederatedResponse>>> readResponses, long[] dims) Bind data from federated workers based on non-overlapping federated ranges.static FederatedRequest[]callInstruction(String[] inst, CPOperand varOldOut, long outputId, CPOperand[] varOldIn, long[] varNewIn, Types.ExecType type) static FederatedRequest[]callInstruction(String[] inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn) static FederatedRequestcallInstruction(String inst, CPOperand varOldOut, long outputId, CPOperand[] varOldIn, long[] varNewIn, Types.ExecType type, boolean rmFedOutputFlag) static FederatedRequestcallInstruction(String inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn) static FederatedRequestcallInstruction(String inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn, boolean rmFedOutFlag) static voidstatic Optional<io.netty.channel.ChannelInboundHandlerAdapter>static Optional<io.netty.channel.ChannelOutboundHandlerAdapter>static Optional<org.apache.commons.lang3.tuple.ImmutablePair<io.netty.channel.ChannelInboundHandlerAdapter,io.netty.channel.ChannelOutboundHandlerAdapter>> static io.netty.handler.codec.serialization.ObjectDecoderdecoder()static FederationMapfederateLocalData(CacheableData<?> data) static longstatic MatrixBlock[]getResults(Future<FederatedResponse>[] ffr) static voidstatic longsumNonZeros(Future<FederatedResponse>[] responses) static voidwaitFor(List<Future<FederatedResponse>> responses)
-
Constructor Details
-
FederationUtils
public FederationUtils()
-
-
Method Details
-
resetFedDataID
public static void resetFedDataID() -
getNextFedDataID
public static long getNextFedDataID() -
checkFedMapType
-
callInstruction
public static FederatedRequest callInstruction(String inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn, boolean rmFedOutFlag) -
callInstruction
public static FederatedRequest callInstruction(String inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn) -
callInstruction
public static FederatedRequest[] callInstruction(String[] inst, CPOperand varOldOut, CPOperand[] varOldIn, long[] varNewIn) -
callInstruction
public static FederatedRequest[] callInstruction(String[] inst, CPOperand varOldOut, long outputId, CPOperand[] varOldIn, long[] varNewIn, Types.ExecType type) -
callInstruction
public static FederatedRequest callInstruction(String inst, CPOperand varOldOut, long outputId, CPOperand[] varOldIn, long[] varNewIn, Types.ExecType type, boolean rmFedOutputFlag) -
aggAdd
-
aggMean
-
getResults
-
bind
-
aggMinMax
public static MatrixBlock aggMinMax(Future<FederatedResponse>[] ffr, boolean isMin, boolean isScalar, Optional<FTypes.FType> fedType) -
aggProd
public static MatrixBlock aggProd(Future<FederatedResponse>[] ffr, FederationMap fedMap, AggregateUnaryOperator aop) -
aggMinMaxIndex
public static MatrixBlock aggMinMaxIndex(Future<FederatedResponse>[] ffr, boolean isMin, FederationMap map) -
aggVar
public static MatrixBlock aggVar(Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map, boolean isRowAggregate, boolean isScalar) -
aggScalar
public static ScalarObject aggScalar(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map) -
aggMatrix
public static MatrixBlock aggMatrix(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, Future<FederatedResponse>[] meanFfr, FederationMap map) -
waitFor
-
aggScalar
-
aggScalar
public static ScalarObject aggScalar(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, FederationMap map) -
aggBooleanScalar
-
aggMatrix
public static MatrixBlock aggMatrix(AggregateUnaryOperator aop, Future<FederatedResponse>[] ffr, FederationMap map) -
federateLocalData
-
bindResponses
public static MatrixBlock bindResponses(List<org.apache.commons.lang3.tuple.Pair<FederatedRange, Future<FederatedResponse>>> readResponses, long[] dims) throws ExceptionBind data from federated workers based on non-overlapping federated ranges.- Parameters:
readResponses- responses from federated workers containing the federated ranges and datadims- dimensions of output MatrixBlock- Returns:
- MatrixBlock of consolidated data
- Throws:
Exception- in case of problems with getting data from responses
-
aggregateResponses
public static MatrixBlock aggregateResponses(List<org.apache.commons.lang3.tuple.Pair<FederatedRange, Future<FederatedResponse>>> readResponses) Aggregate partially aggregated data from federated workers by adding values with the same index in different federated locations.- Parameters:
readResponses- responses from federated workers containing the federated data- Returns:
- MatrixBlock of consolidated, aggregated data
-
decoder
public static io.netty.handler.codec.serialization.ObjectDecoder decoder() -
compressionEncoder
-
compressionDecoder
-
compressionStrategy
public static Optional<org.apache.commons.lang3.tuple.ImmutablePair<io.netty.channel.ChannelInboundHandlerAdapter,io.netty.channel.ChannelOutboundHandlerAdapter>> compressionStrategy() -
sumNonZeros
-