mirror of https://github.com/apache/cassandra
Split out outgoing stream throughput within a DC and inter-DC
patch by Vijay and benedict for CASSANDRA-6596
This commit is contained in:
parent
ee020c9446
commit
94bda1f7d3
|
|
@ -518,6 +518,12 @@ compaction_preheat_key_cache: true
|
||||||
# When unset, the default is 200 Mbps or 25 MB/s.
|
# When unset, the default is 200 Mbps or 25 MB/s.
|
||||||
# stream_throughput_outbound_megabits_per_sec: 200
|
# stream_throughput_outbound_megabits_per_sec: 200
|
||||||
|
|
||||||
|
# Throttles all streaming file transfer between the datacenters,
|
||||||
|
# this setting allows users to throttle inter dc stream throughput in addition
|
||||||
|
# to throttling all network stream traffic as configured with
|
||||||
|
# stream_throughput_outbound_megabits_per_sec
|
||||||
|
# inter_dc_stream_throughput_outbound_megabits_per_sec:
|
||||||
|
|
||||||
# How long the coordinator should wait for read operations to complete
|
# How long the coordinator should wait for read operations to complete
|
||||||
read_request_timeout_in_ms: 5000
|
read_request_timeout_in_ms: 5000
|
||||||
# How long the coordinator should wait for seq or index scans to complete
|
# How long the coordinator should wait for seq or index scans to complete
|
||||||
|
|
|
||||||
|
|
@ -124,6 +124,7 @@ public class Config
|
||||||
public Integer max_streaming_retries = 3;
|
public Integer max_streaming_retries = 3;
|
||||||
|
|
||||||
public volatile Integer stream_throughput_outbound_megabits_per_sec = 200;
|
public volatile Integer stream_throughput_outbound_megabits_per_sec = 200;
|
||||||
|
public volatile Integer inter_dc_stream_throughput_outbound_megabits_per_sec = 0;
|
||||||
|
|
||||||
public String[] data_file_directories;
|
public String[] data_file_directories;
|
||||||
public String flush_directory;
|
public String flush_directory;
|
||||||
|
|
|
||||||
|
|
@ -930,6 +930,16 @@ public class DatabaseDescriptor
|
||||||
conf.stream_throughput_outbound_megabits_per_sec = value;
|
conf.stream_throughput_outbound_megabits_per_sec = value;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static int getInterDCStreamThroughputOutboundMegabitsPerSec()
|
||||||
|
{
|
||||||
|
return conf.inter_dc_stream_throughput_outbound_megabits_per_sec;
|
||||||
|
}
|
||||||
|
|
||||||
|
public static void setInterDCStreamThroughputOutboundMegabitsPerSec(int value)
|
||||||
|
{
|
||||||
|
conf.inter_dc_stream_throughput_outbound_megabits_per_sec = value;
|
||||||
|
}
|
||||||
|
|
||||||
public static String[] getAllDataFileLocations()
|
public static String[] getAllDataFileLocations()
|
||||||
{
|
{
|
||||||
return conf.data_file_directories;
|
return conf.data_file_directories;
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
*/
|
*/
|
||||||
package org.apache.cassandra.streaming;
|
package org.apache.cassandra.streaming;
|
||||||
|
|
||||||
|
import java.net.InetAddress;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
|
|
@ -32,8 +33,8 @@ import com.google.common.collect.Iterables;
|
||||||
import com.google.common.collect.Sets;
|
import com.google.common.collect.Sets;
|
||||||
import com.google.common.util.concurrent.MoreExecutors;
|
import com.google.common.util.concurrent.MoreExecutors;
|
||||||
import com.google.common.util.concurrent.RateLimiter;
|
import com.google.common.util.concurrent.RateLimiter;
|
||||||
import org.cliffc.high_scale_lib.NonBlockingHashMap;
|
|
||||||
|
|
||||||
|
import org.cliffc.high_scale_lib.NonBlockingHashMap;
|
||||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||||
import org.apache.cassandra.streaming.management.StreamEventJMXNotifier;
|
import org.apache.cassandra.streaming.management.StreamEventJMXNotifier;
|
||||||
import org.apache.cassandra.streaming.management.StreamStateCompositeData;
|
import org.apache.cassandra.streaming.management.StreamStateCompositeData;
|
||||||
|
|
@ -47,25 +48,53 @@ public class StreamManager implements StreamManagerMBean
|
||||||
{
|
{
|
||||||
public static final StreamManager instance = new StreamManager();
|
public static final StreamManager instance = new StreamManager();
|
||||||
|
|
||||||
private static final RateLimiter limiter = RateLimiter.create(Double.MAX_VALUE);
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Gets streaming rate limiter.
|
* Gets streaming rate limiter.
|
||||||
* When stream_throughput_outbound_megabits_per_sec is 0, this returns rate limiter
|
* When stream_throughput_outbound_megabits_per_sec is 0, this returns rate limiter
|
||||||
* with the rate of Double.MAX_VALUE bytes per second.
|
* with the rate of Double.MAX_VALUE bytes per second.
|
||||||
* Rate unit is bytes per sec.
|
* Rate unit is bytes per sec.
|
||||||
*
|
*
|
||||||
* @return RateLimiter with rate limit set
|
* @return StreamRateLimiter with rate limit set based on peer location.
|
||||||
*/
|
*/
|
||||||
public static RateLimiter getRateLimiter()
|
public static StreamRateLimiter getRateLimiter(InetAddress peer)
|
||||||
{
|
{
|
||||||
double currentThroughput = ((double) DatabaseDescriptor.getStreamThroughputOutboundMegabitsPerSec()) * 1024 * 1024;
|
return new StreamRateLimiter(peer);
|
||||||
// if throughput is set to 0, throttling is disabled
|
}
|
||||||
if (currentThroughput == 0)
|
|
||||||
currentThroughput = Double.MAX_VALUE;
|
public static class StreamRateLimiter
|
||||||
if (limiter.getRate() != currentThroughput)
|
{
|
||||||
limiter.setRate(currentThroughput);
|
private static final double ONE_MEGA_BITS = 1024 * 1024 * 8;
|
||||||
return limiter;
|
private static final RateLimiter limiter = RateLimiter.create(Double.MAX_VALUE);
|
||||||
|
private static final RateLimiter interDCLimiter = RateLimiter.create(Double.MAX_VALUE);
|
||||||
|
private final boolean isLocalDC;
|
||||||
|
|
||||||
|
public StreamRateLimiter(InetAddress peer)
|
||||||
|
{
|
||||||
|
double throughput = ((double) DatabaseDescriptor.getStreamThroughputOutboundMegabitsPerSec()) * ONE_MEGA_BITS;
|
||||||
|
mayUpdateThroughput(throughput, limiter);
|
||||||
|
|
||||||
|
double interDCThroughput = ((double) DatabaseDescriptor.getInterDCStreamThroughputOutboundMegabitsPerSec()) * ONE_MEGA_BITS;
|
||||||
|
mayUpdateThroughput(interDCThroughput, interDCLimiter);
|
||||||
|
|
||||||
|
isLocalDC = DatabaseDescriptor.getLocalDataCenter().equals(
|
||||||
|
DatabaseDescriptor.getEndpointSnitch().getDatacenter(peer));
|
||||||
|
}
|
||||||
|
|
||||||
|
private void mayUpdateThroughput(double limit, RateLimiter rateLimiter)
|
||||||
|
{
|
||||||
|
// if throughput is set to 0, throttling is disabled
|
||||||
|
if (limit == 0)
|
||||||
|
limit = Double.MAX_VALUE;
|
||||||
|
if (rateLimiter.getRate() != limit)
|
||||||
|
rateLimiter.setRate(limit);
|
||||||
|
}
|
||||||
|
|
||||||
|
public void acquire(int toTransfer)
|
||||||
|
{
|
||||||
|
limiter.acquire(toTransfer);
|
||||||
|
if (!isLocalDC)
|
||||||
|
interDCLimiter.acquire(toTransfer);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private final StreamEventJMXNotifier notifier = new StreamEventJMXNotifier();
|
private final StreamEventJMXNotifier notifier = new StreamEventJMXNotifier();
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,6 @@ import java.nio.channels.Channels;
|
||||||
import java.nio.channels.WritableByteChannel;
|
import java.nio.channels.WritableByteChannel;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
|
|
||||||
import com.google.common.util.concurrent.RateLimiter;
|
|
||||||
import com.ning.compress.lzf.LZFOutputStream;
|
import com.ning.compress.lzf.LZFOutputStream;
|
||||||
|
|
||||||
import org.apache.cassandra.io.sstable.Component;
|
import org.apache.cassandra.io.sstable.Component;
|
||||||
|
|
@ -33,6 +32,7 @@ import org.apache.cassandra.io.util.DataIntegrityMetadata;
|
||||||
import org.apache.cassandra.io.util.DataIntegrityMetadata.ChecksumValidator;
|
import org.apache.cassandra.io.util.DataIntegrityMetadata.ChecksumValidator;
|
||||||
import org.apache.cassandra.io.util.FileUtils;
|
import org.apache.cassandra.io.util.FileUtils;
|
||||||
import org.apache.cassandra.io.util.RandomAccessReader;
|
import org.apache.cassandra.io.util.RandomAccessReader;
|
||||||
|
import org.apache.cassandra.streaming.StreamManager.StreamRateLimiter;
|
||||||
import org.apache.cassandra.utils.Pair;
|
import org.apache.cassandra.utils.Pair;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
@ -44,7 +44,7 @@ public class StreamWriter
|
||||||
|
|
||||||
protected final SSTableReader sstable;
|
protected final SSTableReader sstable;
|
||||||
protected final Collection<Pair<Long, Long>> sections;
|
protected final Collection<Pair<Long, Long>> sections;
|
||||||
protected final RateLimiter limiter = StreamManager.getRateLimiter();
|
protected final StreamRateLimiter limiter;
|
||||||
protected final StreamSession session;
|
protected final StreamSession session;
|
||||||
|
|
||||||
private OutputStream compressedOutput;
|
private OutputStream compressedOutput;
|
||||||
|
|
@ -57,6 +57,7 @@ public class StreamWriter
|
||||||
this.session = session;
|
this.session = session;
|
||||||
this.sstable = sstable;
|
this.sstable = sstable;
|
||||||
this.sections = sections;
|
this.sections = sections;
|
||||||
|
this.limiter = StreamManager.getRateLimiter(session.peer);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue