From b47127b02346efebde43331927945d9300def57d Mon Sep 17 00:00:00 2001 From: Hai Lu Date: Tue, 24 Sep 2019 14:58:05 -0700 Subject: [PATCH 1/3] fix java doc of table descriptor for default rate limiter --- .../samza/table/descriptors/RemoteTableDescriptor.java | 8 ++++++-- .../org/apache/samza/util/EmbeddedTaggedRateLimiter.java | 2 ++ 2 files changed, 8 insertions(+), 2 deletions(-) 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..4821aed25c 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,6 +115,8 @@ public void init(Context context) { int numTasks = jobModel.getContainers().values().stream() .mapToInt(cm -> cm.getTasks().size()) .sum(); + Preconditions.checkArgument(e.getValue() >= numTasks, + String.format("rate limit count (%d) must be greater than number of tasks (%d)", e.getValue(), numTasks)); int effectiveRate = 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, From 0a8f4609295cb1d46d513d4c5331edd168c34a00 Mon Sep 17 00:00:00 2001 From: Hai Lu Date: Tue, 24 Sep 2019 17:10:51 -0700 Subject: [PATCH 2/3] use double in the rate limit calculation --- .../apache/samza/util/EmbeddedTaggedRateLimiter.java | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) 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 4821aed25c..67aff7ee99 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,12 +115,15 @@ public void init(Context context) { int numTasks = jobModel.getContainers().values().stream() .mapToInt(cm -> cm.getTasks().size()) .sum(); - Preconditions.checkArgument(e.getValue() >= numTasks, - String.format("rate limit count (%d) must be greater than number of tasks (%d)", e.getValue(), numTasks)); - 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)) From 24060007884d89b232fb6f5f1c5df66dd5db0c9c Mon Sep 17 00:00:00 2001 From: Hai Lu Date: Tue, 24 Sep 2019 18:46:58 -0700 Subject: [PATCH 3/3] fix style --- .../java/org/apache/samza/util/EmbeddedTaggedRateLimiter.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 67aff7ee99..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 @@ -120,8 +120,8 @@ public void init(Context context) { 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.", + 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));