Class NettyFileTransferClient
java.lang.Object
com.linkedin.davinci.blobtransfer.client.NettyFileTransferClient
-
Field Summary
Fields -
Constructor Summary
ConstructorsConstructorDescriptionNettyFileTransferClient(int serverPort, String baseDir, StorageMetadataService storageMetadataService, int peersConnectivityFreshnessInSeconds, int blobReceiveTimeoutInMin, int blobReceiveReaderIdleTimeInSeconds, int nettyWorkerThreadCount, io.netty.handler.traffic.GlobalChannelTrafficShapingHandler globalChannelTrafficShapingHandler, AggBlobTransferStats aggBlobTransferStats, Optional<SSLFactory> sslFactory, Supplier<VeniceNotifier> notifierSupplier, LogContext logContext, boolean dedicatedAllocatorEnabled) -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()get(String host, String storeName, int version, int partition, BlobTransferUtils.BlobTransferTableFormat requestedTableFormat) io.netty.channel.ChannelgetActiveChannel(String replicaId) Get the active channel for a given transfer, if any.io.netty.buffer.PooledByteBufAllocatorThe allocator the client channels allocate from, which isPooledByteBufAllocator.DEFAULTunless a dedicated one is enabled.getConnectableHosts(HashSet<String> discoveredHosts, String storeName, int version, int partition) A method to get the connectable hosts for the given store, version, and partition This method is only used for checking connectivity to the hosts.intHow many replicas currently hold a transfer channel.voidpurgeStaleConnectivityRecords(VeniceConcurrentHashMap<String, Long> hostsToTimestamp) Check the freshness of the connectivity records and purge the stale records
-
Field Details
-
MIN_NETTY_WORKER_THREADS
public static final int MIN_NETTY_WORKER_THREADS- See Also:
-
-
Constructor Details
-
NettyFileTransferClient
public NettyFileTransferClient(int serverPort, String baseDir, StorageMetadataService storageMetadataService, int peersConnectivityFreshnessInSeconds, int blobReceiveTimeoutInMin, int blobReceiveReaderIdleTimeInSeconds, int nettyWorkerThreadCount, io.netty.handler.traffic.GlobalChannelTrafficShapingHandler globalChannelTrafficShapingHandler, AggBlobTransferStats aggBlobTransferStats, Optional<SSLFactory> sslFactory, Supplier<VeniceNotifier> notifierSupplier, LogContext logContext, boolean dedicatedAllocatorEnabled)
-
-
Method Details
-
getByteBufAllocator
public io.netty.buffer.PooledByteBufAllocator getByteBufAllocator()The allocator the client channels allocate from, which isPooledByteBufAllocator.DEFAULTunless a dedicated one is enabled. -
getConnectableHosts
public Set<String> getConnectableHosts(HashSet<String> discoveredHosts, String storeName, int version, int partition) A method to get the connectable hosts for the given store, version, and partition This method is only used for checking connectivity to the hosts. Channel is closed after checking.- Parameters:
discoveredHosts- the list of discovered hosts for the store, version, and partition, but not necessarily connectablestoreName- the store nameversion- the versionpartition- the partition- Returns:
- the list of connectable hosts
-
purgeStaleConnectivityRecords
Check the freshness of the connectivity records and purge the stale records- Parameters:
hostsToTimestamp- the map of hosts to the timestamp of the last connection attempt
-
get
public CompletionStage<InputStream> get(String host, String storeName, int version, int partition, BlobTransferUtils.BlobTransferTableFormat requestedTableFormat) -
getActiveChannel
Get the active channel for a given transfer, if any.- Parameters:
replicaId- the replica ID (format: storeName_vVersion-partition)- Returns:
- the active Channel, or null if no active transfer
-
getInFlightTransferCount
public int getInFlightTransferCount()How many replicas currently hold a transfer channel. Unlike the fetch executor's active count, this stays elevated for as long as bytes are streaming: a fetch worker hands off to a Netty event loop and returns within milliseconds, while the channel lives until the transfer ends. -
close
public void close()
-