public class DefaultBlockWorkerClient extends java.lang.Object implements BlockWorkerClient
BlockWorkerClient.BlockWorkerClient.Factory| Constructor and Description |
|---|
DefaultBlockWorkerClient(UserState userState,
GrpcServerAddress address,
AlluxioConfiguration alluxioConf)
Creates a client instance for communicating with block worker.
|
| Modifier and Type | Method and Description |
|---|---|
void |
cache(CacheRequest request)
Caches a block.
|
ClearMetricsResponse |
clearMetrics(ClearMetricsRequest request)
Clear the worker metrics.
|
void |
close() |
io.grpc.stub.StreamObserver<CreateLocalBlockRequest> |
createLocalBlock(io.grpc.stub.StreamObserver<CreateLocalBlockResponse> responseObserver)
Creates a local block on the worker.
|
boolean |
isHealthy() |
boolean |
isShutdown() |
com.google.common.util.concurrent.ListenableFuture<LoadResponse> |
load(LoadRequest request)
load blocks into alluxio.
|
MoveBlockResponse |
moveBlock(MoveBlockRequest request)
Move a block from worker.
|
io.grpc.stub.StreamObserver<OpenLocalBlockRequest> |
openLocalBlock(io.grpc.stub.StreamObserver<OpenLocalBlockResponse> responseObserver)
Opens a local block.
|
io.grpc.stub.StreamObserver<ReadRequest> |
readBlock(io.grpc.stub.StreamObserver<ReadResponse> responseObserver)
Reads a block from the worker.
|
RemoveBlockResponse |
removeBlock(RemoveBlockRequest request)
Removes a block from worker.
|
io.grpc.stub.StreamObserver<WriteRequest> |
writeBlock(io.grpc.stub.StreamObserver<WriteResponse> responseObserver)
Writes a block to the worker asynchronously.
|
public DefaultBlockWorkerClient(UserState userState, GrpcServerAddress address, AlluxioConfiguration alluxioConf) throws java.io.IOException
userState - the user stateaddress - the address of the workeralluxioConf - Alluxio configurationjava.io.IOExceptionpublic boolean isShutdown()
isShutdown in interface BlockWorkerClientpublic boolean isHealthy()
isHealthy in interface BlockWorkerClientpublic void close()
throws java.io.IOException
close in interface java.io.Closeableclose in interface java.lang.AutoCloseablejava.io.IOExceptionpublic io.grpc.stub.StreamObserver<WriteRequest> writeBlock(io.grpc.stub.StreamObserver<WriteResponse> responseObserver)
BlockWorkerClientwriteBlock in interface BlockWorkerClientresponseObserver - the stream observer for the server responsepublic io.grpc.stub.StreamObserver<ReadRequest> readBlock(io.grpc.stub.StreamObserver<ReadResponse> responseObserver)
BlockWorkerClientreadBlock in interface BlockWorkerClientresponseObserver - the stream observer for the server responsepublic io.grpc.stub.StreamObserver<CreateLocalBlockRequest> createLocalBlock(io.grpc.stub.StreamObserver<CreateLocalBlockResponse> responseObserver)
BlockWorkerClientcreateLocalBlock in interface BlockWorkerClientresponseObserver - the stream observer for the server responsepublic io.grpc.stub.StreamObserver<OpenLocalBlockRequest> openLocalBlock(io.grpc.stub.StreamObserver<OpenLocalBlockResponse> responseObserver)
BlockWorkerClientopenLocalBlock in interface BlockWorkerClientresponseObserver - the stream observer for the server responsepublic RemoveBlockResponse removeBlock(RemoveBlockRequest request)
BlockWorkerClientremoveBlock in interface BlockWorkerClientrequest - the remove block requestpublic MoveBlockResponse moveBlock(MoveBlockRequest request)
BlockWorkerClientmoveBlock in interface BlockWorkerClientrequest - the remove block requestpublic ClearMetricsResponse clearMetrics(ClearMetricsRequest request)
BlockWorkerClientclearMetrics in interface BlockWorkerClientrequest - the request to clear metricspublic void cache(CacheRequest request)
BlockWorkerClientcache in interface BlockWorkerClientrequest - the cache requestpublic com.google.common.util.concurrent.ListenableFuture<LoadResponse> load(LoadRequest request)
BlockWorkerClientload in interface BlockWorkerClientrequest - the cache requestCopyright © 2022. All Rights Reserved.