diff --git a/samza-api/src/main/java/org/apache/samza/table/descriptors/RemoteTableDescriptor.java b/samza-api/src/main/java/org/apache/samza/table/descriptors/RemoteTableDescriptor.java index 3eed914b6a..23234cdcdf 100644 --- a/samza-api/src/main/java/org/apache/samza/table/descriptors/RemoteTableDescriptor.java +++ b/samza-api/src/main/java/org/apache/samza/table/descriptors/RemoteTableDescriptor.java @@ -225,7 +225,9 @@ public RemoteTableDescriptor 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 withReadRateLimit(int creditsPerSec) { @@ -239,7 +241,9 @@ public RemoteTableDescriptor 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 withWriteRateLimit(int creditsPerSec) { diff --git a/samza-core/src/main/java/org/apache/samza/util/EmbeddedTaggedRateLimiter.java b/samza-core/src/main/java/org/apache/samza/util/EmbeddedTaggedRateLimiter.java index 2bbbf8daef..adb637ecdb 100644 --- a/samza-core/src/main/java/org/apache/samza/util/EmbeddedTaggedRateLimiter.java +++ b/samza-core/src/main/java/org/apache/samza/util/EmbeddedTaggedRateLimiter.java @@ -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))