From 9a7724aab91159da6d0d0187dc7cf98090cd86df Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Mon, 16 Sep 2019 16:07:56 -0700 Subject: [PATCH 1/9] Adding logic to read system config for repl-factor when creating a topic --- .../apache/samza/system/kafka/KafkaSystemAdmin.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java index c18c82dc2a..44c391e644 100644 --- a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java +++ b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java @@ -546,6 +546,19 @@ public KafkaStreamSpec toKafkaSpec(StreamSpec spec) { kafkaSpec = kafkaSpec.copyWithProperties(properties); } else { kafkaSpec = KafkaStreamSpec.fromSpec(spec); + + // we check if there is a system-level rf config specified, else we use KafkaConfig.topic-default-rf + int replicationFactorFromSystemConfig = Integer.valueOf(new SystemConfig(config). + getDefaultStreamProperties(spec.getSystemName()).getOrDefault(KafkaConfig.TOPIC_REPLICATION_FACTOR(), + KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR())); + + LOG.info("Using replication-factor: {} for StreamSpec: {}", replicationFactorFromSystemConfig, spec); + + return new KafkaStreamSpec( kafkaSpec.getId(), + kafkaSpec.getPhysicalName(), + kafkaSpec.getSystemName(), + kafkaSpec.getPartitionCount(), + replicationFactorFromSystemConfig,kafkaSpec.getProperties()); } return kafkaSpec; } From 7006ba71a5f7ce99153d76e804fa9752cec541eb Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Mon, 16 Sep 2019 16:52:36 -0700 Subject: [PATCH 2/9] Adding fix for changelog system --- .../org/apache/samza/config/KafkaConfig.scala | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index b9543fbaab..5fbdcc9c94 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -243,17 +243,21 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { /** * Gets the replication factor for the changelog topics. Uses the following precedence. * - * 1. If stores.myStore.changelog.replication.factor is configured, that value is used. - * 2. If systems.changelog-system.default.stream.replication.factor is configured, that value is used. - * 3. 2 - * - * Note that the changelog-system has a similar precedence. See [[StorageConfig]] + * 1. If the changelog system (specified using job.changelog.system or job.default.system is provided), the RF for that system is used. + * 2. If neither of these are provided, the RF specified using stores..changelog.replication.factor is used. + * 3. If this is not provided, the RF specified using stores.default.changelog.replication.factor is used. + * 4. If this is not provided, the RF is chosen as 2. */ - def getChangelogStreamReplicationFactor(name: String) = getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format name).getOrElse(getDefaultChangelogStreamReplicationFactor) - - def getDefaultChangelogStreamReplicationFactor() = { + def getChangelogStreamReplicationFactor(name: String) = { val changelogSystem = new StorageConfig(config).getChangelogSystem.orElse(null) - getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR).getOrElse(getSystemDefaultReplicationFactor(changelogSystem, "2")) + val systemDefaultReplicationFactor = getSystemDefaultReplicationFactor(changelogSystem, null) + + if (systemDefaultReplicationFactor == null) { + getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format name).getOrElse( + getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR).getOrElse("2")) + } else { + systemDefaultReplicationFactor + } } /** From 1375945f14c383ba795663925c096101da96c93c Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 17 Sep 2019 12:00:42 -0700 Subject: [PATCH 3/9] Updating getChangelogStream(name) in StorageConfig --- .../apache/samza/config/StorageConfig.java | 2 +- .../samza/config/TestStorageConfig.java | 47 ++++++++++++++----- .../org/apache/samza/config/KafkaConfig.scala | 39 ++++++++++----- .../apache/samza/config/TestKafkaConfig.scala | 3 -- 4 files changed, 61 insertions(+), 30 deletions(-) diff --git a/samza-core/src/main/java/org/apache/samza/config/StorageConfig.java b/samza-core/src/main/java/org/apache/samza/config/StorageConfig.java index 7bc6cb4c72..86c7e7de3b 100644 --- a/samza-core/src/main/java/org/apache/samza/config/StorageConfig.java +++ b/samza-core/src/main/java/org/apache/samza/config/StorageConfig.java @@ -148,7 +148,7 @@ public Optional getStorageMsgSerde(String storeName) { * * @return the name of the system to use by default for all changelogs, if defined. */ - public Optional getChangelogSystem() { + private Optional getChangelogSystem() { return Optional.ofNullable(get(CHANGELOG_SYSTEM, get(JobConfig.JOB_DEFAULT_SYSTEM))); } diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java index 2cde5df560..63931dd0cc 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java @@ -26,8 +26,10 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import org.apache.samza.SamzaException; +import org.apache.samza.util.StreamUtil; import org.junit.Test; +import static org.apache.samza.config.StorageConfig.*; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -59,24 +61,24 @@ public void testGetChangelogStream() { // store has empty string for changelog stream StorageConfig storageConfig = new StorageConfig( - new MapConfig(ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), ""))); + new MapConfig(ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), ""))); assertEquals(Optional.empty(), storageConfig.getChangelogStream(STORE_NAME0)); // store has full changelog system-stream defined storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), + ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.changelog-stream0"))); assertEquals(Optional.of("changelog-system.changelog-stream0"), storageConfig.getChangelogStream(STORE_NAME0)); // store has changelog stream defined, but system comes from job.changelog.system storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0", + ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0", StorageConfig.CHANGELOG_SYSTEM, "changelog-system"))); assertEquals(Optional.of("changelog-system.changelog-stream0"), storageConfig.getChangelogStream(STORE_NAME0)); // batch mode: create unique stream name storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), + ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.changelog-stream0", ApplicationConfig.APP_MODE, ApplicationConfig.ApplicationMode.BATCH.name().toLowerCase(), ApplicationConfig.APP_RUN_ID, "run-id"))); assertEquals(Optional.of("changelog-system.changelog-stream0-run-id"), @@ -86,7 +88,7 @@ public void testGetChangelogStream() { @Test(expected = SamzaException.class) public void testGetChangelogStreamMissingSystem() { StorageConfig storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0"))); + ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0"))); storageConfig.getChangelogStream(STORE_NAME0); } @@ -155,18 +157,37 @@ public void testGetStorageMsgSerde() { @Test public void testGetChangelogSystem() { // empty config, so no system - assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogSystem()); + assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogStream(STORE_NAME0)); - // job.changelog.system takes precedence over job.default.system StorageConfig storageConfig = new StorageConfig(new MapConfig( ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, "should-not-be-used"))); - assertEquals(Optional.of("changelog-system"), storageConfig.getChangelogSystem()); + assertEquals(Optional.empty(), storageConfig.getChangelogStream(STORE_NAME0)); + + // job.changelog.system takes precedence over job.default.system when changelog is specified as just streamName + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, + "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); + assertEquals("changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + + // job.changelog.system takes precedence over job.default.system when changelog is specified as {systemName}.{streamName} + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, + "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.streamName"))); + assertEquals("changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + + // systemName specified using stores.{storeName}.changelog = {systemName}.{streamName} should take precedence even + // when job.changelog.system and job.default.system are specified + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "default-changelog-system", + JobConfig.JOB_DEFAULT_SYSTEM, "default-system", + String.format(CHANGELOG_STREAM, STORE_NAME0), "nondefault-changelog-system.streamName"))); + assertEquals("nondefault-changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); // fall back to job.default.system if job.changelog.system is not specified - storageConfig = - new StorageConfig(new MapConfig(ImmutableMap.of(JobConfig.JOB_DEFAULT_SYSTEM, "default-system"))); - assertEquals(Optional.of("default-system"), storageConfig.getChangelogSystem()); + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(JobConfig.JOB_DEFAULT_SYSTEM, "default-system", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); + assertEquals("default-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); } @Test @@ -232,7 +253,7 @@ public void testIsChangelogSystem() { StorageConfig storageConfig = new StorageConfig(new MapConfig(ImmutableMap.of( // store0 has a changelog stream String.format(StorageConfig.FACTORY, STORE_NAME0), "factory.class", - String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "system0.changelog-stream", + String.format(CHANGELOG_STREAM, STORE_NAME0), "system0.changelog-stream", // store1 does not have a changelog stream String.format(StorageConfig.FACTORY, STORE_NAME1), "factory.class"))); assertTrue(storageConfig.isChangelogSystem("system0")); @@ -248,7 +269,7 @@ public void testHasDurableStores() { storageConfig = new StorageConfig(new MapConfig( ImmutableMap.of(String.format(StorageConfig.FACTORY, STORE_NAME0), "factory.class", - String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "system0.changelog-stream"))); + String.format(CHANGELOG_STREAM, STORE_NAME0), "system0.changelog-stream"))); assertTrue(storageConfig.hasDurableStores()); } diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 5fbdcc9c94..44e2e29acc 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -30,6 +30,7 @@ import org.apache.kafka.clients.producer.ProducerConfig import org.apache.kafka.common.serialization.ByteArraySerializer import org.apache.samza.SamzaException import org.apache.samza.config.ApplicationConfig.ApplicationMode +import org.apache.samza.system.SystemStream import org.apache.samza.util.ScalaJavaUtil.JavaOptionals import org.apache.samza.util.{Logging, StreamUtil} @@ -243,23 +244,35 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { /** * Gets the replication factor for the changelog topics. Uses the following precedence. * - * 1. If the changelog system (specified using job.changelog.system or job.default.system is provided), the RF for that system is used. - * 2. If neither of these are provided, the RF specified using stores..changelog.replication.factor is used. - * 3. If this is not provided, the RF specified using stores.default.changelog.replication.factor is used. - * 4. If this is not provided, the RF is chosen as 2. + * 1. If stores.{storeName}.changelog.replication.factor is configured, that value is used. + * 2. If it is not configured, the value configured for stores.default.changelog.replication.factor is used. + * 3. If it is not configured, the RF value configured for the store's changelog's system, configured using + * stores.{storeName}.changelog={systemName}.{streamName}, is used. + * 4. If it is not configured, the value for the RF of job.changelog.system is used. + * 5. If it is not configured, the value for the RF of job.default.system is used. + * 6. If it is not configured, the RF is chosen as 2. */ - def getChangelogStreamReplicationFactor(name: String) = { - val changelogSystem = new StorageConfig(config).getChangelogSystem.orElse(null) - val systemDefaultReplicationFactor = getSystemDefaultReplicationFactor(changelogSystem, null) + def getChangelogStreamReplicationFactor(storeName: String) = { + var changelogRF = getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format storeName) - if (systemDefaultReplicationFactor == null) { - getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format name).getOrElse( - getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR).getOrElse("2")) - } else { - systemDefaultReplicationFactor + if(!changelogRF.isDefined) { + changelogRF = getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR) + } + + if(!changelogRF.isDefined) { + val changelogSystemStream = new StorageConfig(config).getChangelogStream(storeName) + if (!changelogSystemStream.isPresent) { + throw new SamzaException("Changelog system not defined for store "+storeName) + } + + val changelogSystem = StreamUtil.getSystemStreamFromNames(changelogSystemStream.get()).getSystem + changelogRF = Option.apply(getSystemDefaultReplicationFactor(changelogSystem, "2")) } + + changelogRF.get } + /** * Gets the max message bytes for the changelog topics. Uses the following precedence. * @@ -272,7 +285,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { def getChangelogStreamMaxMessageByte(name: String) = getOption(KafkaConfig.CHANGELOG_MAX_MESSAGE_BYTES format name) match { case Some(maxMessageBytes) => maxMessageBytes case _ => - val changelogSystem = new StorageConfig(config).getChangelogSystem.orElse(null) + val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).get()).getSystem val systemMaxMessageBytes = new SystemConfig(config).getDefaultStreamProperties(changelogSystem).getOrDefault(KafkaConfig.MAX_MESSAGE_BYTES, KafkaConfig.DEFAULT_LOG_COMPACT_TOPIC_MAX_MESSAGE_BYTES) systemMaxMessageBytes } diff --git a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala index bb1b337326..083919beaf 100644 --- a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala +++ b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala @@ -220,7 +220,6 @@ class TestKafkaConfig { val kafkaConfig = new KafkaConfig(mapConfig) assertEquals("3", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override")) assertEquals("2", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override")) - assertEquals("2", kafkaConfig.getDefaultChangelogStreamReplicationFactor) } @Test @@ -235,7 +234,6 @@ class TestKafkaConfig { val kafkaConfig = new KafkaConfig(mapConfig) assertEquals("4", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override")) assertEquals("5", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override")) - assertEquals("5", kafkaConfig.getDefaultChangelogStreamReplicationFactor) } @Test @@ -248,7 +246,6 @@ class TestKafkaConfig { val kafkaConfig = new KafkaConfig(mapConfig) assertEquals("4", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override")) assertEquals("8", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override")) - assertEquals("8", kafkaConfig.getDefaultChangelogStreamReplicationFactor) } @Test From b79c98e6674cd84d4db3f73691690fc83a91a8ae Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 17 Sep 2019 12:03:37 -0700 Subject: [PATCH 4/9] Undoing superfluous change --- .../org/apache/samza/config/TestStorageConfig.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java index 63931dd0cc..2ec3de82f1 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java @@ -61,24 +61,24 @@ public void testGetChangelogStream() { // store has empty string for changelog stream StorageConfig storageConfig = new StorageConfig( - new MapConfig(ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), ""))); + new MapConfig(ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), ""))); assertEquals(Optional.empty(), storageConfig.getChangelogStream(STORE_NAME0)); // store has full changelog system-stream defined storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), + ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-system.changelog-stream0"))); assertEquals(Optional.of("changelog-system.changelog-stream0"), storageConfig.getChangelogStream(STORE_NAME0)); // store has changelog stream defined, but system comes from job.changelog.system storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0", + ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0", StorageConfig.CHANGELOG_SYSTEM, "changelog-system"))); assertEquals(Optional.of("changelog-system.changelog-stream0"), storageConfig.getChangelogStream(STORE_NAME0)); // batch mode: create unique stream name storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), + ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-system.changelog-stream0", ApplicationConfig.APP_MODE, ApplicationConfig.ApplicationMode.BATCH.name().toLowerCase(), ApplicationConfig.APP_RUN_ID, "run-id"))); assertEquals(Optional.of("changelog-system.changelog-stream0-run-id"), @@ -88,7 +88,7 @@ public void testGetChangelogStream() { @Test(expected = SamzaException.class) public void testGetChangelogStreamMissingSystem() { StorageConfig storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0"))); + ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), "changelog-stream0"))); storageConfig.getChangelogStream(STORE_NAME0); } From 86518146756af6deb09082a255153378225509d1 Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 17 Sep 2019 12:14:29 -0700 Subject: [PATCH 5/9] Simplyfying KafkaStreamSpec.toKafkaSpec --- .../org/apache/samza/system/kafka/KafkaSystemAdmin.java | 6 +++--- .../main/scala/org/apache/samza/config/KafkaConfig.scala | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java index 44c391e644..0905c98f7a 100644 --- a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java +++ b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java @@ -548,9 +548,9 @@ public KafkaStreamSpec toKafkaSpec(StreamSpec spec) { kafkaSpec = KafkaStreamSpec.fromSpec(spec); // we check if there is a system-level rf config specified, else we use KafkaConfig.topic-default-rf - int replicationFactorFromSystemConfig = Integer.valueOf(new SystemConfig(config). - getDefaultStreamProperties(spec.getSystemName()).getOrDefault(KafkaConfig.TOPIC_REPLICATION_FACTOR(), - KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR())); + int replicationFactorFromSystemConfig = Integer.valueOf( + new KafkaConfig(config).getSystemDefaultReplicationFactor(spec.getSystemName(), + KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR())); LOG.info("Using replication-factor: {} for StreamSpec: {}", replicationFactorFromSystemConfig, spec); diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 44e2e29acc..7a3c68e98c 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -119,7 +119,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { Option(replicationFactor) } - private def getSystemDefaultReplicationFactor(systemName: String, defaultValue: String) = { + def getSystemDefaultReplicationFactor(systemName: String, defaultValue: String) = { val defaultReplicationFactor = new SystemConfig(config).getDefaultStreamProperties(systemName).getOrDefault(KafkaConfig.TOPIC_REPLICATION_FACTOR, defaultValue) defaultReplicationFactor } From e9fd653cc81a3c439709591e04507b7f99ff3bf5 Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Tue, 17 Sep 2019 15:15:06 -0700 Subject: [PATCH 6/9] Applying review comments, committing tests --- .../org/apache/samza/config/TestStorageConfig.java | 11 +++++------ .../apache/samza/system/kafka/KafkaSystemAdmin.java | 9 ++------- .../scala/org/apache/samza/config/KafkaConfig.scala | 9 ++++----- .../org/apache/samza/config/TestKafkaConfig.scala | 13 +++++++++++++ 4 files changed, 24 insertions(+), 18 deletions(-) diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java index 2ec3de82f1..b724660e3c 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java @@ -26,7 +26,6 @@ import com.google.common.collect.ImmutableMap; import com.google.common.collect.ImmutableSet; import org.apache.samza.SamzaException; -import org.apache.samza.util.StreamUtil; import org.junit.Test; import static org.apache.samza.config.StorageConfig.*; @@ -155,7 +154,7 @@ public void testGetStorageMsgSerde() { } @Test - public void testGetChangelogSystem() { + public void getChangelogStream() { // empty config, so no system assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogStream(STORE_NAME0)); @@ -168,13 +167,13 @@ public void testGetChangelogSystem() { storageConfig = new StorageConfig(new MapConfig( ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); - assertEquals("changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); // job.changelog.system takes precedence over job.default.system when changelog is specified as {systemName}.{streamName} storageConfig = new StorageConfig(new MapConfig( ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.streamName"))); - assertEquals("changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); // systemName specified using stores.{storeName}.changelog = {systemName}.{streamName} should take precedence even // when job.changelog.system and job.default.system are specified @@ -182,12 +181,12 @@ public void testGetChangelogSystem() { ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "default-changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, "default-system", String.format(CHANGELOG_STREAM, STORE_NAME0), "nondefault-changelog-system.streamName"))); - assertEquals("nondefault-changelog-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + assertEquals("nondefault-changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); // fall back to job.default.system if job.changelog.system is not specified storageConfig = new StorageConfig(new MapConfig( ImmutableMap.of(JobConfig.JOB_DEFAULT_SYSTEM, "default-system", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); - assertEquals("default-system", StreamUtil.getSystemStreamFromNames(storageConfig.getChangelogStream(STORE_NAME0).get()).getSystem()); + assertEquals("default-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); } @Test diff --git a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java index 0905c98f7a..1986bea54f 100644 --- a/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java +++ b/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java @@ -551,14 +551,9 @@ public KafkaStreamSpec toKafkaSpec(StreamSpec spec) { int replicationFactorFromSystemConfig = Integer.valueOf( new KafkaConfig(config).getSystemDefaultReplicationFactor(spec.getSystemName(), KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR())); - LOG.info("Using replication-factor: {} for StreamSpec: {}", replicationFactorFromSystemConfig, spec); - - return new KafkaStreamSpec( kafkaSpec.getId(), - kafkaSpec.getPhysicalName(), - kafkaSpec.getSystemName(), - kafkaSpec.getPartitionCount(), - replicationFactorFromSystemConfig,kafkaSpec.getProperties()); + return new KafkaStreamSpec(kafkaSpec.getId(), kafkaSpec.getPhysicalName(), kafkaSpec.getSystemName(), + kafkaSpec.getPartitionCount(), replicationFactorFromSystemConfig, kafkaSpec.getProperties()); } return kafkaSpec; } diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 7a3c68e98c..70837da462 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -30,7 +30,6 @@ import org.apache.kafka.clients.producer.ProducerConfig import org.apache.kafka.common.serialization.ByteArraySerializer import org.apache.samza.SamzaException import org.apache.samza.config.ApplicationConfig.ApplicationMode -import org.apache.samza.system.SystemStream import org.apache.samza.util.ScalaJavaUtil.JavaOptionals import org.apache.samza.util.{Logging, StreamUtil} @@ -262,14 +261,14 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { if(!changelogRF.isDefined) { val changelogSystemStream = new StorageConfig(config).getChangelogStream(storeName) if (!changelogSystemStream.isPresent) { - throw new SamzaException("Changelog system not defined for store "+storeName) + throw new SamzaException("Cannot deduce replication factor. Changelog system-stream not defined for store " + storeName) } val changelogSystem = StreamUtil.getSystemStreamFromNames(changelogSystemStream.get()).getSystem - changelogRF = Option.apply(getSystemDefaultReplicationFactor(changelogSystem, "2")) + changelogRF = Option.apply(getSystemDefaultReplicationFactor(changelogSystem, null)) } - changelogRF.get + changelogRF.getOrElse(KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR) } @@ -285,7 +284,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { def getChangelogStreamMaxMessageByte(name: String) = getOption(KafkaConfig.CHANGELOG_MAX_MESSAGE_BYTES format name) match { case Some(maxMessageBytes) => maxMessageBytes case _ => - val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).get()).getSystem + val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).orElse(null)).getSystem val systemMaxMessageBytes = new SystemConfig(config).getDefaultStreamProperties(changelogSystem).getOrDefault(KafkaConfig.MAX_MESSAGE_BYTES, KafkaConfig.DEFAULT_LOG_COMPACT_TOPIC_MAX_MESSAGE_BYTES) systemMaxMessageBytes } diff --git a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala index 083919beaf..d94f414b32 100644 --- a/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala +++ b/samza-kafka/src/test/scala/org/apache/samza/config/TestKafkaConfig.scala @@ -211,10 +211,21 @@ class TestKafkaConfig { kafkaProducerConfig.getProducerProperties } + @Test + def testGetSystemDefaultReplicationFactor(): Unit = { + assertEquals(KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR, new KafkaConfig(new MapConfig()).getSystemDefaultReplicationFactor("kafka-system",KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR)) + + props.setProperty("systems.kafka-system.default.stream.replication.factor", "8") + val mapConfig = new MapConfig(props.asScala.asJava) + val kafkaConfig = new KafkaConfig(mapConfig) + assertEquals("8", kafkaConfig.getSystemDefaultReplicationFactor("kafka-system","2")) + } + @Test def testChangeLogReplicationFactor() { props.setProperty("stores.store-with-override.changelog", "kafka-system.changelog-topic") props.setProperty("stores.store-with-override.changelog.replication.factor", "3") + props.setProperty("stores.default.changelog.replication.factor", "2") val mapConfig = new MapConfig(props.asScala.asJava) val kafkaConfig = new KafkaConfig(mapConfig) @@ -230,6 +241,7 @@ class TestKafkaConfig { // Override the "default" default value props.setProperty("stores.default.changelog.replication.factor", "5") + val mapConfig = new MapConfig(props.asScala.asJava) val kafkaConfig = new KafkaConfig(mapConfig) assertEquals("4", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override")) @@ -241,6 +253,7 @@ class TestKafkaConfig { props.setProperty(StorageConfig.CHANGELOG_SYSTEM, "kafka-system") props.setProperty("systems.kafka-system.default.stream.replication.factor", "8") props.setProperty("stores.store-with-override.changelog.replication.factor", "4") + props.setProperty("stores.store-without-override.changelog", "change-for-store-without-override") val mapConfig = new MapConfig(props.asScala.asJava) val kafkaConfig = new KafkaConfig(mapConfig) From 2dd52b069440a7d63ba85f7c8b45eee13bba873e Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Wed, 18 Sep 2019 13:51:13 -0700 Subject: [PATCH 7/9] Addressing review comments --- .../samza/config/TestStorageConfig.java | 70 ++++++++----------- .../org/apache/samza/config/KafkaConfig.scala | 2 +- 2 files changed, 32 insertions(+), 40 deletions(-) diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java index b724660e3c..948151f8e9 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java @@ -55,9 +55,6 @@ public void testGetStoreNames() { @Test public void testGetChangelogStream() { - // empty config, so no changelog stream - assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogStream(STORE_NAME0)); - // store has empty string for changelog stream StorageConfig storageConfig = new StorageConfig( new MapConfig(ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), ""))); @@ -82,6 +79,37 @@ public void testGetChangelogStream() { ApplicationConfig.ApplicationMode.BATCH.name().toLowerCase(), ApplicationConfig.APP_RUN_ID, "run-id"))); assertEquals(Optional.of("changelog-system.changelog-stream0-run-id"), storageConfig.getChangelogStream(STORE_NAME0)); + + // job has no changelog stream defined + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, + "should-not-be-used"))); + assertEquals(Optional.empty(), storageConfig.getChangelogStream(STORE_NAME0)); + + // job.changelog.system takes precedence over job.default.system when changelog is specified as just streamName + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, + "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); + assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); + + // job.changelog.system takes precedence over job.default.system when changelog is specified as {systemName}.{streamName} + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, + "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.streamName"))); + assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); + + // systemName specified using stores.{storeName}.changelog = {systemName}.{streamName} should take precedence even + // when job.changelog.system and job.default.system are specified + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "default-changelog-system", + JobConfig.JOB_DEFAULT_SYSTEM, "default-system", + String.format(CHANGELOG_STREAM, STORE_NAME0), "nondefault-changelog-system.streamName"))); + assertEquals("nondefault-changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); + + // fall back to job.default.system if job.changelog.system is not specified + storageConfig = new StorageConfig(new MapConfig( + ImmutableMap.of(JobConfig.JOB_DEFAULT_SYSTEM, "default-system", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); + assertEquals("default-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); } @Test(expected = SamzaException.class) @@ -153,42 +181,6 @@ public void testGetStorageMsgSerde() { assertEquals(Optional.of("my.msg.serde.class"), storageConfig.getStorageMsgSerde(STORE_NAME0)); } - @Test - public void getChangelogStream() { - // empty config, so no system - assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogStream(STORE_NAME0)); - - StorageConfig storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, - "should-not-be-used"))); - assertEquals(Optional.empty(), storageConfig.getChangelogStream(STORE_NAME0)); - - // job.changelog.system takes precedence over job.default.system when changelog is specified as just streamName - storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, - "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); - assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); - - // job.changelog.system takes precedence over job.default.system when changelog is specified as {systemName}.{streamName} - storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "changelog-system", JobConfig.JOB_DEFAULT_SYSTEM, - "should-not-be-used", String.format(CHANGELOG_STREAM, STORE_NAME0), "changelog-system.streamName"))); - assertEquals("changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); - - // systemName specified using stores.{storeName}.changelog = {systemName}.{streamName} should take precedence even - // when job.changelog.system and job.default.system are specified - storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(StorageConfig.CHANGELOG_SYSTEM, "default-changelog-system", - JobConfig.JOB_DEFAULT_SYSTEM, "default-system", - String.format(CHANGELOG_STREAM, STORE_NAME0), "nondefault-changelog-system.streamName"))); - assertEquals("nondefault-changelog-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); - - // fall back to job.default.system if job.changelog.system is not specified - storageConfig = new StorageConfig(new MapConfig( - ImmutableMap.of(JobConfig.JOB_DEFAULT_SYSTEM, "default-system", String.format(CHANGELOG_STREAM, STORE_NAME0), "streamName"))); - assertEquals("default-system.streamName", storageConfig.getChangelogStream(STORE_NAME0).get()); - } - @Test public void testGetSideInputs() { // empty config, so no system diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 70837da462..025cb98439 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -284,7 +284,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { def getChangelogStreamMaxMessageByte(name: String) = getOption(KafkaConfig.CHANGELOG_MAX_MESSAGE_BYTES format name) match { case Some(maxMessageBytes) => maxMessageBytes case _ => - val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).orElse(null)).getSystem + val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).orElseThrow(() => new SamzaException("System-stream not defined for store:"+name))).getSystem val systemMaxMessageBytes = new SystemConfig(config).getDefaultStreamProperties(changelogSystem).getOrDefault(KafkaConfig.MAX_MESSAGE_BYTES, KafkaConfig.DEFAULT_LOG_COMPACT_TOPIC_MAX_MESSAGE_BYTES) systemMaxMessageBytes } From 83213dc1fdabd3dc7e745f739969ccfeddfebabb Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Wed, 18 Sep 2019 14:13:47 -0700 Subject: [PATCH 8/9] Adding test removed in error --- .../test/java/org/apache/samza/config/TestStorageConfig.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java index 948151f8e9..e094de2333 100644 --- a/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java +++ b/samza-core/src/test/java/org/apache/samza/config/TestStorageConfig.java @@ -55,6 +55,9 @@ public void testGetStoreNames() { @Test public void testGetChangelogStream() { + // empty config, so no changelog stream + assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogStream(STORE_NAME0)); + // store has empty string for changelog stream StorageConfig storageConfig = new StorageConfig( new MapConfig(ImmutableMap.of(String.format(StorageConfig.CHANGELOG_STREAM, STORE_NAME0), ""))); From 9d05eca2396ef7e21d4ba39eec5c3cea935fc2a9 Mon Sep 17 00:00:00 2001 From: Ray Matharu Date: Wed, 18 Sep 2019 14:53:06 -0700 Subject: [PATCH 9/9] Updating getOrElse --- .../src/main/scala/org/apache/samza/config/KafkaConfig.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala index 025cb98439..993f0e4225 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala +++ b/samza-kafka/src/main/scala/org/apache/samza/config/KafkaConfig.scala @@ -284,7 +284,7 @@ class KafkaConfig(config: Config) extends ScalaMapConfig(config) { def getChangelogStreamMaxMessageByte(name: String) = getOption(KafkaConfig.CHANGELOG_MAX_MESSAGE_BYTES format name) match { case Some(maxMessageBytes) => maxMessageBytes case _ => - val changelogSystem = StreamUtil.getSystemStreamFromNames(new StorageConfig(config).getChangelogStream(name).orElseThrow(() => new SamzaException("System-stream not defined for store:"+name))).getSystem + val changelogSystem = StreamUtil.getSystemStreamFromNames(JavaOptionals.toRichOptional(new StorageConfig(config).getChangelogStream(name)).toOption.getOrElse(throw new SamzaException("System-stream not defined for store:"+name))).getSystem val systemMaxMessageBytes = new SystemConfig(config).getDefaultStreamProperties(changelogSystem).getOrDefault(KafkaConfig.MAX_MESSAGE_BYTES, KafkaConfig.DEFAULT_LOG_COMPACT_TOPIC_MAX_MESSAGE_BYTES) systemMaxMessageBytes }