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 @@ -226,6 +226,10 @@ public boolean isCoordinatorStream() {
return id.equals(COORDINATOR_STREAM_ID);
}

public boolean isCheckpointStream() {
return id.equals(CHECKPOINT_STREAM_ID);
}

private void validateLogicalIdentifier(String identifierName, String identifierValue) {
if (identifierValue == null || !identifierValue.matches("[A-Za-z0-9_-]+")) {
throw new IllegalArgumentException(String.format("Identifier '%s' is '%s'. It must match the expression [A-Za-z0-9_-]+", identifierName, identifierValue));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -540,6 +540,9 @@ public KafkaStreamSpec toKafkaSpec(StreamSpec spec) {
kafkaSpec =
new KafkaStreamSpec(spec.getId(), spec.getPhysicalName(), systemName, 1, coordinatorStreamReplicationFactor,
coordinatorStreamProperties);
} else if (spec.isCheckpointStream()) {
kafkaSpec = KafkaStreamSpec.fromSpec(StreamSpec.createCheckpointStreamSpec(spec.getPhysicalName(), systemName))
.copyWithReplicationFactor(Integer.parseInt(new KafkaConfig(config).getCheckpointReplicationFactor().get()));
Comment on lines +543 to +545

@rmatharu-zz rmatharu-zz Oct 3, 2019 •

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.

This makes sense, however, KafkaCheckpointManagerFactory already adheres to these configs when creating the spec.
Please see KafkaCheckpointManagerFactory Line 52-55.

val checkpointSpec = KafkaStreamSpec.fromSpec(StreamSpec.createCheckpointStreamSpec(checkpointTopic, checkpointSystemName))
.copyWithReplicationFactor(kafkaConfig.getCheckpointReplicationFactor.get.toInt)
.copyWithProperties(kafkaConfig.getCheckpointTopicProperties)

@VartulAgrawal VartulAgrawal Oct 3, 2019 •

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.

The KafkaCheckpointManagerFactory is only created after the KafkaSystemAdmin is initialised so by the time KafkaCheckpointManagerFactory is called the replicationFactor is already reset to 2 (https://github.com/apache/samza/blob/master/samza-kafka/src/main/java/org/apache/samza/system/kafka/KafkaSystemAdmin.java#L552 which is default for any general topic) and hence it needs to be passed to the spec here. Tested this in our current setup too :)

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.

Thanks. I'll merge your patch for now.
The problem really is that the KafkaSystemAdmin's toStreamSpec does not respect the incoming repl-factor in case the stream-spec is already a KafkaStreamSpec and tries to differentiate based on the stream-ID.

I've summarized the problem here https://issues.apache.org/jira/browse/SAMZA-2342
for a better, clean, long-term fix.

} else if (intermediateStreamProperties.containsKey(spec.getId())) {
kafkaSpec = KafkaStreamSpec.fromSpec(spec);
Properties properties = kafkaSpec.getProperties();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -271,6 +271,23 @@ class TestKafkaSystemAdmin {
validateTopic(topic, 1)
}

@Test
def testShouldCreateCheckpointStream {
val topic = "test-checkpoint-stream"
val expectedReplicationFactor = 1
val map = new java.util.HashMap[String, String]()
map.put(org.apache.samza.config.KafkaConfig.CHECKPOINT_REPLICATION_FACTOR, expectedReplicationFactor.toString)
val systemAdmin = createSystemAdmin(SYSTEM, map)

val spec = StreamSpec.createCheckpointStreamSpec(topic, "kafka")
val actualReplicationFactor = systemAdmin.toKafkaSpec(spec).getReplicationFactor
// Ensure we respect the replication factor passed to system admin
assertEquals(expectedReplicationFactor, actualReplicationFactor)

systemAdmin.createStream(spec)
validateTopic(topic, 1)
}

@Test
def testGetNewestOffset {
createTopic(TOPIC2, 16)
Expand Down