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 @@ -148,7 +148,7 @@ public Optional<String> getStorageMsgSerde(String storeName) {
*
* @return the name of the system to use by default for all changelogs, if defined.
*/
public Optional<String> getChangelogSystem() {
private Optional<String> getChangelogSystem() {
return Optional.ofNullable(get(CHANGELOG_SYSTEM, get(JobConfig.JOB_DEFAULT_SYSTEM)));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.samza.SamzaException;
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;
Expand Down Expand Up @@ -81,6 +82,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)
Expand Down Expand Up @@ -152,23 +184,6 @@ public void testGetStorageMsgSerde() {
assertEquals(Optional.of("my.msg.serde.class"), storageConfig.getStorageMsgSerde(STORE_NAME0));
}

@Test
public void testGetChangelogSystem() {
// empty config, so no system
assertEquals(Optional.empty(), new StorageConfig(new MapConfig()).getChangelogSystem());

// 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());

// 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());
}

@Test
public void testGetSideInputs() {
// empty config, so no system
Expand Down Expand Up @@ -232,7 +247,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"));
Expand All @@ -248,7 +263,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());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -546,6 +546,14 @@ public KafkaStreamSpec toKafkaSpec(StreamSpec spec) {
kafkaSpec = kafkaSpec.copyWithProperties(properties);
} else {
kafkaSpec = KafkaStreamSpec.fromSpec(spec);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

To clarify my understanding: Is RF correctly configurable for the above cases (changelog, coordinator, intermediate streams), but not other cases? So would an example case impacted by this block be diagnostics?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

yes, and for changelog not all configs were being respected (see getChangelogStreamReplicationFactor)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah, yes, right. Although the fix for changelog was done in KafkaConfig, and then the fix for some of the other stream types was done here.
Thanks for clarifying.

// we check if there is a system-level rf config specified, else we use KafkaConfig.topic-default-rf
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 kafkaSpec;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,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
}
Expand Down Expand Up @@ -243,19 +243,35 @@ 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 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) = getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format name).getOrElse(getDefaultChangelogStreamReplicationFactor)
def getChangelogStreamReplicationFactor(storeName: String) = {
var changelogRF = getOption(KafkaConfig.CHANGELOG_STREAM_REPLICATION_FACTOR format storeName)

if(!changelogRF.isDefined) {
changelogRF = getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR)
}

def getDefaultChangelogStreamReplicationFactor() = {
val changelogSystem = new StorageConfig(config).getChangelogSystem.orElse(null)
getOption(KafkaConfig.DEFAULT_CHANGELOG_STREAM_REPLICATION_FACTOR).getOrElse(getSystemDefaultReplicationFactor(changelogSystem, "2"))
if(!changelogRF.isDefined) {
val changelogSystemStream = new StorageConfig(config).getChangelogStream(storeName)
if (!changelogSystemStream.isPresent) {
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, null))
}

changelogRF.getOrElse(KafkaConfig.TOPIC_DEFAULT_REPLICATION_FACTOR)
}


/**
* Gets the max message bytes for the changelog topics. Uses the following precedence.
*
Expand All @@ -268,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 = new StorageConfig(config).getChangelogSystem.orElse(null)
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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -211,16 +211,26 @@ 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)
assertEquals("3", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override"))
assertEquals("2", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override"))
assertEquals("2", kafkaConfig.getDefaultChangelogStreamReplicationFactor)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can you please add unit tests for the new flows in KafkaConfig for RF?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

added in TestKafkaConfig and updated others

}

@Test
Expand All @@ -231,24 +241,24 @@ 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"))
assertEquals("5", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override"))
assertEquals("5", kafkaConfig.getDefaultChangelogStreamReplicationFactor)
}

@Test
def testChangeLogReplicationFactorWithSystemOverriddenDefault() {
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)
assertEquals("4", kafkaConfig.getChangelogStreamReplicationFactor("store-with-override"))
assertEquals("8", kafkaConfig.getChangelogStreamReplicationFactor("store-without-override"))
assertEquals("8", kafkaConfig.getDefaultChangelogStreamReplicationFactor)
}

@Test
Expand Down