Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,9 @@ public RemoteTableDescriptor<K, V> withWriteRateLimiterDisabled() {
* it is invalid to call {@link RemoteTableDescriptor#withRateLimiter(RateLimiter,
* TableRateLimiter.CreditFunction, TableRateLimiter.CreditFunction)}
* and vice versa.
* @param creditsPerSec rate limit for read operations; must be positive
* Note that this is the total credit of rate limit for the entire job, each task will get a per task
* credit of creditsPerSec/tasksCount. Hence creditsPerSec should be greater than total number of tasks.
* @param creditsPerSec rate limit for read operations; must be positive and greater than total number tasks
* @return this table descriptor instance
*/
public RemoteTableDescriptor<K, V> withReadRateLimit(int creditsPerSec) {
Expand All @@ -239,7 +241,9 @@ public RemoteTableDescriptor<K, V> withReadRateLimit(int creditsPerSec) {
* it is invalid to call {@link RemoteTableDescriptor#withRateLimiter(RateLimiter,
* TableRateLimiter.CreditFunction, TableRateLimiter.CreditFunction)}
* and vice versa.
* @param creditsPerSec rate limit for write operations; must be positive
* Note that this is the total credit of rate limit for the entire job, each task will get a per task
* credit of creditsPerSec/tasksCount. Hence creditsPerSec should be greater than total number of tasks.
* @param creditsPerSec rate limit for write operations; must be positive and greater than total number tasks
* @return this table descriptor instance
*/
public RemoteTableDescriptor<K, V> withWriteRateLimit(int creditsPerSec) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -115,10 +115,15 @@ public void init(Context context) {
int numTasks = jobModel.getContainers().values().stream()
.mapToInt(cm -> cm.getTasks().size())
.sum();
int effectiveRate = e.getValue() / numTasks;
double effectiveRate = (double) e.getValue() / numTasks;
TaskName taskName = context.getTaskContext().getTaskModel().getTaskName();
LOGGER.info(String.format("Effective rate limit for task %s and tag %s is %d", taskName, tag,
LOGGER.info(String.format("Effective rate limit for task %s and tag %s is %f", taskName, tag,
effectiveRate));
if (effectiveRate < 1.0) {
LOGGER.warn(String.format("Effective limit rate (%f) is very low. "
+ "Total rate limit is %d while number of tasks is %d. Consider increasing the rate limit.",
effectiveRate, e.getValue(), numTasks));
}
return new ImmutablePair<>(tag, com.google.common.util.concurrent.RateLimiter.create(effectiveRate));
})
.collect(Collectors.toMap(ImmutablePair::getKey, ImmutablePair::getValue))
Expand Down