From d5aeab92967501832ec4de2288c437711f02f692 Mon Sep 17 00:00:00 2001 From: vartulagrawal Date: Thu, 3 Oct 2019 14:25:13 +0100 Subject: [PATCH] Update KafkaSystemAdmin to adhere to the checkpoint replication factor --- .../org/apache/samza/system/StreamSpec.java | 4 ++++ .../samza/system/kafka/KafkaSystemAdmin.java | 3 +++ .../system/kafka/TestKafkaSystemAdmin.scala | 17 +++++++++++++++++ 3 files changed, 24 insertions(+) diff --git a/samza-api/src/main/java/org/apache/samza/system/StreamSpec.java b/samza-api/src/main/java/org/apache/samza/system/StreamSpec.java index aa71f0ed86..a1ad5e4cde 100644 --- a/samza-api/src/main/java/org/apache/samza/system/StreamSpec.java +++ b/samza-api/src/main/java/org/apache/samza/system/StreamSpec.java @@ -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)); 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 f0bce193d9..97229db9eb 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 @@ -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())); } else if (intermediateStreamProperties.containsKey(spec.getId())) { kafkaSpec = KafkaStreamSpec.fromSpec(spec); Properties properties = kafkaSpec.getProperties(); diff --git a/samza-kafka/src/test/scala/org/apache/samza/system/kafka/TestKafkaSystemAdmin.scala b/samza-kafka/src/test/scala/org/apache/samza/system/kafka/TestKafkaSystemAdmin.scala index 9dceb5e968..9bd8acfa8f 100644 --- a/samza-kafka/src/test/scala/org/apache/samza/system/kafka/TestKafkaSystemAdmin.scala +++ b/samza-kafka/src/test/scala/org/apache/samza/system/kafka/TestKafkaSystemAdmin.scala @@ -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)