From 7cc1bda2f3ffa90793eac19420b4d99d61ebe1d7 Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Fri, 13 Sep 2019 15:59:25 -0700 Subject: [PATCH 1/6] SAMZA-2322: Integrate startpoints with ZK Standalone integration tests. --- .../test/processor/TestStreamApplication.java | 5 +- .../test/processor/TestTaskApplication.java | 13 +- .../TestZkLocalApplicationRunner.java | 352 ++++++++++++++---- 3 files changed, 303 insertions(+), 67 deletions(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java index af6afc12b8..923291cc23 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java @@ -24,8 +24,9 @@ import java.io.Serializable; import java.util.List; import java.util.concurrent.CountDownLatch; -import org.apache.samza.application.descriptors.StreamApplicationDescriptor; +import org.apache.commons.lang3.StringUtils; import org.apache.samza.application.StreamApplication; +import org.apache.samza.application.descriptors.StreamApplicationDescriptor; import org.apache.samza.config.ApplicationConfig; import org.apache.samza.config.Config; import org.apache.samza.config.JobConfig; @@ -141,7 +142,7 @@ public String toString() { } static TestKafkaEvent fromString(String message) { - String[] messageComponents = message.split("|"); + String[] messageComponents = StringUtils.split(message, "|"); return new TestKafkaEvent(messageComponents[0], messageComponents[1]); } } diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java index c1723b645f..4a635cc1aa 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java @@ -46,14 +46,16 @@ public class TestTaskApplication implements TaskApplication { private final String outputTopic; private final CountDownLatch shutdownLatch; private final CountDownLatch processedMessageLatch; + private final TaskApplicationCallback processCallback; public TestTaskApplication(String systemName, String inputTopic, String outputTopic, - CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch) { + CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch, TaskApplicationCallback processCallback) { this.systemName = systemName; this.inputTopic = inputTopic; this.outputTopic = outputTopic; this.processedMessageLatch = processedMessageLatch; this.shutdownLatch = shutdownLatch; + this.processCallback = processCallback; } private class TestTaskImpl implements AsyncStreamTask, ClosableTask { @@ -62,7 +64,10 @@ private class TestTaskImpl implements AsyncStreamTask, ClosableTask { public void processAsync(IncomingMessageEnvelope envelope, MessageCollector collector, TaskCoordinator coordinator, TaskCallback callback) { processedMessageLatch.countDown(); // Implementation does not invokes callback.complete to block the RunLoop.process() after it exhausts the - // `task.max.concurrency` defined per task. + // `task.max.concurrency` defined per task. Call callback.complete() in the processCallback if needed. + if (processCallback != null) { + processCallback.onMessage(envelope, callback); + } } @Override @@ -81,4 +86,8 @@ public void describe(TaskApplicationDescriptor appDescriptor) { .withOutputStream(outputDescriptor) .withTaskFactory((AsyncStreamTaskFactory) () -> new TestTaskImpl()); } + + public interface TaskApplicationCallback { + void onMessage(IncomingMessageEnvelope m, TaskCallback callback); + } } diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java index 192243c54f..ae095ad177 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java @@ -24,32 +24,38 @@ import com.google.common.collect.ImmutableSet; import com.google.common.collect.Maps; import com.google.common.collect.Sets; +import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; -import java.util.UUID; import java.util.Objects; import java.util.Set; -import java.util.HashSet; +import java.util.TreeMap; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; import org.I0Itec.zkclient.ZkClient; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.samza.Partition; +import org.apache.samza.SamzaException; import org.apache.samza.application.TaskApplication; import org.apache.samza.config.ApplicationConfig; -import org.apache.samza.config.Config; import org.apache.samza.config.ClusterManagerConfig; +import org.apache.samza.config.Config; import org.apache.samza.config.JobConfig; import org.apache.samza.config.JobCoordinatorConfig; import org.apache.samza.config.MapConfig; import org.apache.samza.config.TaskConfig; import org.apache.samza.config.ZkConfig; -import org.apache.samza.SamzaException; import org.apache.samza.container.TaskName; import org.apache.samza.coordinator.metadatastore.CoordinatorStreamStore; import org.apache.samza.coordinator.stream.CoordinatorStreamValueSerde; @@ -65,15 +71,27 @@ import org.apache.samza.runtime.ApplicationRunner; import org.apache.samza.runtime.ApplicationRunners; import org.apache.samza.runtime.LocalApplicationRunner; +import org.apache.samza.startpoint.Startpoint; +import org.apache.samza.startpoint.StartpointManager; +import org.apache.samza.startpoint.StartpointOldest; +import org.apache.samza.startpoint.StartpointSpecific; +import org.apache.samza.startpoint.StartpointTimestamp; +import org.apache.samza.startpoint.StartpointUpcoming; +import org.apache.samza.system.IncomingMessageEnvelope; +import org.apache.samza.system.SystemAdmin; +import org.apache.samza.system.SystemAdmins; +import org.apache.samza.system.SystemStream; import org.apache.samza.system.SystemStreamPartition; +import org.apache.samza.task.TaskCallback; import org.apache.samza.test.StandaloneTestUtils; import org.apache.samza.test.harness.IntegrationTestHarness; +import org.apache.samza.util.CoordinatorStreamUtil; import org.apache.samza.util.NoOpMetricsRegistry; import org.apache.samza.util.ReflectionUtil; -import org.apache.samza.zk.ZkMetadataStore; -import org.apache.samza.zk.ZkStringSerializer; import org.apache.samza.zk.ZkJobCoordinatorFactory; import org.apache.samza.zk.ZkKeyBuilder; +import org.apache.samza.zk.ZkMetadataStore; +import org.apache.samza.zk.ZkStringSerializer; import org.apache.samza.zk.ZkUtils; import org.junit.Assert; import org.junit.Rule; @@ -83,7 +101,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertTrue; /** @@ -108,9 +129,12 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private static final String TASK_SHUTDOWN_MS = "10000"; private static final String JOB_DEBOUNCE_TIME_MS = "10000"; private static final String BARRIER_TIMEOUT_MS = "10000"; - private static final String[] PROCESSOR_IDS = new String[] {"0000000000", "0000000001", "0000000002"}; + private static final String[] PROCESSOR_IDS = new String[] {"0000000000", "0000000001", "0000000002", "0000000003"}; - private String inputKafkaTopic; + private String inputKafkaTopic1; + private String inputKafkaTopic2; + private String inputKafkaTopic3; + private String inputKafkaTopic4; private String outputKafkaTopic; private String inputSinglePartitionKafkaTopic; private String outputSinglePartitionKafkaTopic; @@ -118,6 +142,7 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private ApplicationConfig applicationConfig1; private ApplicationConfig applicationConfig2; private ApplicationConfig applicationConfig3; + private ApplicationConfig applicationConfig4; private String testStreamAppName; private String testStreamAppId; private MetadataStore zkMetadataStore; @@ -134,7 +159,10 @@ public void setUp() { String uniqueTestId = UUID.randomUUID().toString(); testStreamAppName = String.format("test-app-name-%s", uniqueTestId); testStreamAppId = String.format("test-app-id-%s", uniqueTestId); - inputKafkaTopic = String.format("test-input-topic-%s", uniqueTestId); + inputKafkaTopic1 = String.format("test-input-topic1-%s", uniqueTestId); + inputKafkaTopic2 = String.format("test-input-topic2-%s", uniqueTestId); + inputKafkaTopic3 = String.format("test-input-topic3-%s", uniqueTestId); + inputKafkaTopic4 = String.format("test-input-topic4-%s", uniqueTestId); outputKafkaTopic = String.format("test-output-topic-%s", uniqueTestId); inputSinglePartitionKafkaTopic = String.format("test-input-single-partition-topic-%s", uniqueTestId); outputSinglePartitionKafkaTopic = String.format("test-output-single-partition-topic-%s", uniqueTestId); @@ -149,20 +177,26 @@ public void setUp() { applicationConfig2 = new ApplicationConfig(new MapConfig(configMap)); configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[2]); applicationConfig3 = new ApplicationConfig(new MapConfig(configMap)); + configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[3]); + applicationConfig4 = new ApplicationConfig(new MapConfig(configMap)); ZkClient zkClient = new ZkClient(zkConnect(), ZK_CONNECTION_TIMEOUT_MS, ZK_CONNECTION_TIMEOUT_MS, new ZkStringSerializer()); ZkKeyBuilder zkKeyBuilder = new ZkKeyBuilder(ZkJobCoordinatorFactory.getJobCoordinationZkPath(applicationConfig1)); zkUtils = new ZkUtils(zkKeyBuilder, zkClient, ZK_CONNECTION_TIMEOUT_MS, ZK_SESSION_TIMEOUT_MS, new NoOpMetricsRegistry()); zkUtils.connect(); - ImmutableMap topicToPartitionCount = ImmutableMap.of( - inputSinglePartitionKafkaTopic, 1, - outputSinglePartitionKafkaTopic, 1, - inputKafkaTopic, 5, - outputKafkaTopic, 5); + ImmutableMap topicToPartitionCount = ImmutableMap.builder() + .put(inputSinglePartitionKafkaTopic, 1) + .put(outputSinglePartitionKafkaTopic, 1) + .put(inputKafkaTopic1, 5) + .put(inputKafkaTopic2, 5) + .put(inputKafkaTopic3, 5) + .put(inputKafkaTopic4, 5) + .put(outputKafkaTopic, 5) + .build(); List newTopics = - ImmutableList.of(inputKafkaTopic, outputKafkaTopic, inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic) + topicToPartitionCount.keySet() .stream() .map(topic -> new NewTopic(topic, topicToPartitionCount.get(topic), (short) 1)) .collect(Collectors.toList()); @@ -173,8 +207,8 @@ public void setUp() { @Override public void tearDown() { - deleteTopics(ImmutableList.of( - inputKafkaTopic, outputKafkaTopic, inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic)); + deleteTopics(ImmutableList.of(inputKafkaTopic1, inputKafkaTopic2, inputKafkaTopic3, inputKafkaTopic4, outputKafkaTopic, + inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic)); SharedContextFactories.clearAll(); zkUtils.close(); super.tearDown(); @@ -184,12 +218,36 @@ private void publishKafkaEvents(String topic, int startIndex, int endIndex, Stri for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { try { LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); - producer.send(new ProducerRecord(topic, new TestStreamApplication.TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)).toString().getBytes())); + TestStreamApplication.TestKafkaEvent testKafkaEvent = + new TestStreamApplication.TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)); + producer.send(new ProducerRecord(topic, testKafkaEvent.toString().getBytes())); + } catch (Exception e) { + LOGGER.error("Publishing to kafka topic: {} resulted in exception: {}.", new Object[]{topic, e}); + throw new SamzaException(e); + } + } + producer.flush(); + } + + // Sends each event synchronously and returns map of event index to event metadata of the sent record/event. Adds the specified delay between sends. + private TreeMap publishKafkaEventsWithDelayPerEvent(String topic, int startIndex, int endIndex, String streamProcessorId, Duration delay) { + TreeMap eventsMetadata = new TreeMap<>(); + for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { + try { + LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); + TestStreamApplication.TestKafkaEvent testKafkaEvent = + new TestStreamApplication.TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)); + Thread.sleep(delay.toMillis()); + Future send = producer.send(new ProducerRecord(topic, + testKafkaEvent.toString().getBytes())); + eventsMetadata.put(eventIndex, send.get(5, TimeUnit.SECONDS)); } catch (Exception e) { LOGGER.error("Publishing to kafka topic: {} resulted in exception: {}.", new Object[]{topic, e}); throw new SamzaException(e); } } + producer.flush(); + return eventsMetadata; } private Map buildStreamApplicationConfigMap(String appName, String appId, boolean isBatch) { @@ -327,7 +385,7 @@ public void shouldStopNewProcessorsJoiningGroupWhenNumContainersIsGreaterThanNum @Test public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleTaskGrouperFactory() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS * 2, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS * 2, PROCESSOR_IDS[0]); // Configuration, verification variables MapConfig testConfig = new MapConfig(ImmutableMap.of(JobConfig.SSP_GROUPER_FACTORY, @@ -349,7 +407,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT CountDownLatch processedMessagesLatch = new CountDownLatch(NUM_KAFKA_EVENTS * 2); Config testAppConfig2 = new MapConfig(applicationConfig2, testConfig); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch, null, null, testAppConfig2), testAppConfig2); // Callback handler for appRunner1. @@ -373,7 +431,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT // Set up stream app appRunner1. Config testAppConfig1 = new MapConfig(applicationConfig1, testConfig); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, null, streamApplicationCallback, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, null, streamApplicationCallback, kafkaEventsConsumedLatch, testAppConfig1), testAppConfig1); executeRun(appRunner1, testAppConfig1); @@ -418,7 +476,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT @Test public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create stream applications. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(2 * NUM_KAFKA_EVENTS); @@ -427,13 +485,13 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, applicationConfig3), applicationConfig3); executeRun(appRunner1, applicationConfig1); @@ -461,7 +519,7 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { assertEquals(ApplicationStatus.SuccessfulFinish, appRunner1.status()); kafkaEventsConsumedLatch.await(); - publishKafkaEvents(inputKafkaTopic, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -486,7 +544,7 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { @Test public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create StreamApplications. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); @@ -494,10 +552,10 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Run stream applications. @@ -509,10 +567,10 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw processedMessagesLatch2.await(); // Create a stream app with same processor id as SP2 and run it. It should fail. - publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[2]); + publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[2]); kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, null, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, null, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Fail when the duplicate processor joins. expectedException.expect(SamzaException.class); @@ -527,10 +585,14 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw } } + public void addMessagesProcessed(List messagesProcessed, TestStreamApplication.TestKafkaEvent tke) { + messagesProcessed.add(tke); + } + @Test public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() throws Exception { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId); @@ -541,7 +603,8 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t Config applicationConfig2 = new MapConfig(configMap); List messagesProcessed = new ArrayList<>(); - TestStreamApplication.StreamApplicationCallback streamApplicationCallback = messagesProcessed::add; + TestStreamApplication.StreamApplicationCallback streamApplicationCallback = + (TestStreamApplication.TestKafkaEvent tke) -> addMessagesProcessed(messagesProcessed, tke); // Create StreamApplication from configuration. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); @@ -549,10 +612,10 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, streamApplicationCallback, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Run stream application. @@ -578,9 +641,9 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t assertEquals(ApplicationStatus.SuccessfulFinish, appRunner1.status()); processedMessagesLatch1 = new CountDownLatch(1); - publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); executeRun(appRunner3, applicationConfig1); @@ -603,7 +666,7 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t @Test public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContainerShutdownTime() throws Exception { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId); configMap.put(TaskConfig.TASK_SHUTDOWN_MS, "0"); @@ -620,10 +683,10 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -645,11 +708,11 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); - publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); processedMessagesLatch3.await(); appRunner1.waitForFinish(); @@ -675,14 +738,14 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine */ @Test public void testShouldGenerateJobModelOnPartitionCountChange() throws Exception { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create StreamApplication from configuration. CountDownLatch kafkaEventsConsumedLatch1 = new CountDownLatch(NUM_KAFKA_EVENTS); CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch1, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch1, applicationConfig1), applicationConfig1); executeRun(appRunner1, applicationConfig1); @@ -696,7 +759,7 @@ public void testShouldGenerateJobModelOnPartitionCountChange() throws Exception Assert.assertEquals(5, ssps.size()); // Increase the partition count of input kafka topic to 100. - increasePartitionsTo(inputKafkaTopic, 100); + increasePartitionsTo(inputKafkaTopic1, 100); long jobModelWaitTimeInMillis = 10; while (Objects.equals(zkUtils.getJobModelVersion(), jobModelVersion)) { @@ -938,14 +1001,14 @@ public void testStatefulSamzaApplicationShouldRedistributeInputPartitionsToCorre @Test public void testApplicationShutdownShouldBeIndependentOfPerMessageProcessingTime() throws Exception { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create a TaskApplication with only one task per container. // The task does not invokes taskCallback.complete for any of the dispatched message. CountDownLatch shutdownLatch = new CountDownLatch(1); CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); - TaskApplication taskApplication = new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic, outputKafkaTopic, processedMessagesLatch1, shutdownLatch); + TaskApplication taskApplication = new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatch1, shutdownLatch, null); MapConfig taskApplicationConfig = new MapConfig(ImmutableList.of(applicationConfig1, ImmutableMap.of(TaskConfig.MAX_CONCURRENCY, "1", JobConfig.SSP_GROUPER_FACTORY, "org.apache.samza.container.grouper.stream.AllSspToSingleTaskGrouperFactory"))); ApplicationRunner appRunner = ApplicationRunners.getApplicationRunner(taskApplication, taskApplicationConfig); @@ -975,7 +1038,7 @@ public void testApplicationShutdownShouldBeIndependentOfPerMessageProcessingTime */ @Test public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -990,10 +1053,10 @@ public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1027,7 +1090,7 @@ public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { */ @Test public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1043,10 +1106,10 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1066,7 +1129,7 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1097,7 +1160,7 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException */ @Test public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1113,10 +1176,10 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1144,7 +1207,7 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1171,7 +1234,7 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr */ @Test public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1184,7 +1247,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); executeRun(appRunner1, applicationConfig1); @@ -1196,7 +1259,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep // bring up second processor CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner2, applicationConfig2); @@ -1213,7 +1276,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1230,6 +1293,159 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep appRunner3.waitForFinish(); } + private void writeStartpoints(StartpointManager startpointManager, String streamName, Integer partitionCount, Startpoint startpoint) { + for (int p = 0; p < partitionCount; p++) { + startpointManager.writeStartpoint(new SystemStreamPartition(TEST_SYSTEM, streamName, new Partition(p)), startpoint); + } + } + + @Test + public void testStartpoints() throws InterruptedException { + TreeMap sentEvents1 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0], Duration.ofMillis(2)); + TreeMap sentEvents2 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic2, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[1], Duration.ofMillis(2)); + TreeMap sentEvents3 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic3, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[2], Duration.ofMillis(2)); + TreeMap sentEvents4 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic4, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[3], Duration.ofMillis(2)); + + ConcurrentHashMap recvEventsInputStartpointSpecific = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointTimestamp = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointOldest = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointUpcoming = new ConcurrentHashMap<>(); + + CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); + coordinatorStreamStore.init(); + StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); + startpointManager.start(); + + StartpointSpecific startpointSpecific = new StartpointSpecific(String.valueOf(sentEvents1.get(100).offset())); + writeStartpoints(startpointManager, inputKafkaTopic1, 5, startpointSpecific); + + StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); + writeStartpoints(startpointManager, inputKafkaTopic2, 5, startpointTimestamp); + + StartpointOldest startpointOldest = new StartpointOldest(); + writeStartpoints(startpointManager, inputKafkaTopic3, 5, startpointOldest); + + StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); + writeStartpoints(startpointManager, inputKafkaTopic4, 5, startpointUpcoming); + + startpointManager.stop(); + coordinatorStreamStore.close(); + + TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + try { + String streamName = ime.getSystemStreamPartition().getStream(); + TestStreamApplication.TestKafkaEvent testKafkaEvent = + TestStreamApplication.TestKafkaEvent.fromString((String) ime.getMessage()); + String eventIndex = testKafkaEvent.getEventData(); + if (inputKafkaTopic1.equals(streamName)) { + recvEventsInputStartpointSpecific.put(eventIndex, ime); + } else if (inputKafkaTopic2.equals(streamName)) { + recvEventsInputStartpointTimestamp.put(eventIndex, ime); + } else if (inputKafkaTopic3.equals(streamName)) { + recvEventsInputStartpointOldest.put(eventIndex, ime); + } else if (inputKafkaTopic4.equals(streamName)) { + recvEventsInputStartpointUpcoming.put(eventIndex, ime); + } else { + throw new RuntimeException("Unexpected input stream: " + streamName); + } + callback.complete(); + } catch (Exception ex) { + callback.failure(ex); + } + }; + + // Create StreamApplication from configuration. + CountDownLatch processedMessagesLatchStartpointSpecific = new CountDownLatch(100); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Just fetch a few messages + + CountDownLatch shutdownLatchStartpointSpecific = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointOldest = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointUpcoming = new CountDownLatch(1); + + TestTaskApplication testTaskApplicationStartpointSpecific = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatchStartpointSpecific, + shutdownLatchStartpointSpecific, processedCallback); + TestTaskApplication testTaskApplicationStartpointTimestamp = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic2, outputKafkaTopic, processedMessagesLatchStartpointTimestamp, + shutdownLatchStartpointTimestamp, processedCallback); + TestTaskApplication testTaskApplicationStartpointOldest = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic3, outputKafkaTopic, processedMessagesLatchStartpointOldest, + shutdownLatchStartpointOldest, processedCallback); + TestTaskApplication testTaskApplicationStartpointUpcoming = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic4, outputKafkaTopic, processedMessagesLatchStartpointUpcoming, + shutdownLatchStartpointUpcoming, processedCallback); + + // Startpoint specific + ApplicationRunner appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointSpecific, applicationConfig1); + executeRun(appRunner, applicationConfig1); + + assertTrue(processedMessagesLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); + + Integer startpointSpecificOffset = Integer.valueOf(startpointSpecific.getSpecificOffset()); + for (IncomingMessageEnvelope ime : recvEventsInputStartpointSpecific.values()) { + Integer eventOffset = Integer.valueOf(ime.getOffset()); + String assertMsg = String.format("Expecting message offset: %d >= Startpoint specific offset: %d", + eventOffset, startpointSpecificOffset); + assertTrue(assertMsg, eventOffset >= startpointSpecificOffset); + } + + // Startpoint timestamp + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointTimestamp, applicationConfig2); + + executeRun(appRunner, applicationConfig2); + assertTrue(processedMessagesLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); + + for (IncomingMessageEnvelope ime : recvEventsInputStartpointTimestamp.values()) { + Integer eventOffset = Integer.valueOf(ime.getOffset()); + assertNotEquals(0, ime.getEventTime()); // sanity check + String assertMsg = String.format("Expecting message timestamp: %d >= Startpoint timestamp: %d", + ime.getEventTime(), startpointTimestamp.getTimestampOffset()); + assertTrue(assertMsg, ime.getEventTime() >= startpointTimestamp.getTimestampOffset()); + } + + // Startpoint oldest + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointOldest, applicationConfig3); + executeRun(appRunner, applicationConfig3); + + assertTrue(processedMessagesLatchStartpointOldest.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointOldest.await(1, TimeUnit.MINUTES)); + + assertEquals("Expecting to have processed all the events", NUM_KAFKA_EVENTS, recvEventsInputStartpointOldest.size()); + + // Startpoint upcoming + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointUpcoming, applicationConfig4); + executeRun(appRunner, applicationConfig4); + + assertFalse("Expecting to timeout and not process any old messages.", processedMessagesLatchStartpointUpcoming.await(15, TimeUnit.SECONDS)); + assertEquals("Expecting not to process any old messages.", 0, recvEventsInputStartpointUpcoming.size()); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointUpcoming.await(1, TimeUnit.MINUTES)); + } + /** * Computes the task to partition assignment of the {@param JobModel}. * @param jobModel the jobModel to compute task to partition assignment for. @@ -1258,4 +1474,14 @@ private static List getSystemStreamPartitions(JobModel jo }); return ssps; } + + private static CoordinatorStreamStore createCoordinatorStreamStore(Config applicationConfig) { + SystemStream coordinatorSystemStream = CoordinatorStreamUtil.getCoordinatorSystemStream(applicationConfig); + SystemAdmins systemAdmins = new SystemAdmins(applicationConfig); + SystemAdmin coordinatorSystemAdmin = systemAdmins.getSystemAdmin(coordinatorSystemStream.getSystem()); + coordinatorSystemAdmin.start(); + CoordinatorStreamUtil.createCoordinatorStream(coordinatorSystemStream, coordinatorSystemAdmin); + coordinatorSystemAdmin.stop(); + return new CoordinatorStreamStore(applicationConfig, new NoOpMetricsRegistry()); + } } From 6f578922480c3108aaf8a684424d334fcd042085 Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Mon, 16 Sep 2019 09:22:28 -0700 Subject: [PATCH 2/6] Minor cleanup. --- .../TestZkLocalApplicationRunner.java | 23 ++++++++++--------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java index ae095ad177..01a1279b52 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java @@ -121,6 +121,7 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private static final int NUM_KAFKA_EVENTS = 300; private static final int ZK_CONNECTION_TIMEOUT_MS = 5000; private static final int ZK_SESSION_TIMEOUT_MS = 10000; + private static final int ZK_TEST_PARTITION_COUNT = 5; private static final String TEST_SYSTEM = "TestSystemName"; private static final String TEST_SSP_GROUPER_FACTORY = "org.apache.samza.container.grouper.stream.GroupByPartitionFactory"; private static final String TEST_TASK_GROUPER_FACTORY = "org.apache.samza.container.grouper.task.GroupByContainerIdsFactory"; @@ -188,11 +189,11 @@ public void setUp() { ImmutableMap topicToPartitionCount = ImmutableMap.builder() .put(inputSinglePartitionKafkaTopic, 1) .put(outputSinglePartitionKafkaTopic, 1) - .put(inputKafkaTopic1, 5) - .put(inputKafkaTopic2, 5) - .put(inputKafkaTopic3, 5) - .put(inputKafkaTopic4, 5) - .put(outputKafkaTopic, 5) + .put(inputKafkaTopic1, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic2, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic3, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic4, ZK_TEST_PARTITION_COUNT) + .put(outputKafkaTopic, ZK_TEST_PARTITION_COUNT) .build(); List newTopics = @@ -1321,16 +1322,16 @@ public void testStartpoints() throws InterruptedException { startpointManager.start(); StartpointSpecific startpointSpecific = new StartpointSpecific(String.valueOf(sentEvents1.get(100).offset())); - writeStartpoints(startpointManager, inputKafkaTopic1, 5, startpointSpecific); + writeStartpoints(startpointManager, inputKafkaTopic1, ZK_TEST_PARTITION_COUNT, startpointSpecific); StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); - writeStartpoints(startpointManager, inputKafkaTopic2, 5, startpointTimestamp); + writeStartpoints(startpointManager, inputKafkaTopic2, ZK_TEST_PARTITION_COUNT, startpointTimestamp); StartpointOldest startpointOldest = new StartpointOldest(); - writeStartpoints(startpointManager, inputKafkaTopic3, 5, startpointOldest); + writeStartpoints(startpointManager, inputKafkaTopic3, ZK_TEST_PARTITION_COUNT, startpointOldest); StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); - writeStartpoints(startpointManager, inputKafkaTopic4, 5, startpointUpcoming); + writeStartpoints(startpointManager, inputKafkaTopic4, ZK_TEST_PARTITION_COUNT, startpointUpcoming); startpointManager.stop(); coordinatorStreamStore.close(); @@ -1361,8 +1362,8 @@ public void testStartpoints() throws InterruptedException { // Create StreamApplication from configuration. CountDownLatch processedMessagesLatchStartpointSpecific = new CountDownLatch(100); // Just fetch a few messages CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Fetch all since consuming from oldest + CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Expecting none, so just attempt a small number of fetches. CountDownLatch shutdownLatchStartpointSpecific = new CountDownLatch(1); CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); From 45c98fcfef752ca018b75a9db303964494caa024 Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Mon, 23 Sep 2019 15:14:10 -0700 Subject: [PATCH 3/6] Move startpoint test from TestZkLocalApplicationRunner into its own test harness class. --- .../test/processor/TestStreamApplication.java | 35 +- .../test/processor/TestTaskApplication.java | 31 +- .../TestZkLocalApplicationRunner.java | 341 +++------------- .../startpoint/StartpointTestHarness.java | 383 ++++++++++++++++++ .../samza/test/util/TestKafkaEvent.java | 58 +++ 5 files changed, 527 insertions(+), 321 deletions(-) create mode 100644 samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java create mode 100644 samza-test/src/test/java/org/apache/samza/test/util/TestKafkaEvent.java diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java index 923291cc23..4c3bc4598b 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestStreamApplication.java @@ -21,10 +21,8 @@ import java.io.IOException; import java.io.ObjectInputStream; -import java.io.Serializable; import java.util.List; import java.util.concurrent.CountDownLatch; -import org.apache.commons.lang3.StringUtils; import org.apache.samza.application.StreamApplication; import org.apache.samza.application.descriptors.StreamApplicationDescriptor; import org.apache.samza.config.ApplicationConfig; @@ -38,6 +36,7 @@ import org.apache.samza.system.kafka.descriptors.KafkaInputDescriptor; import org.apache.samza.system.kafka.descriptors.KafkaOutputDescriptor; import org.apache.samza.system.kafka.descriptors.KafkaSystemDescriptor; +import org.apache.samza.test.util.TestKafkaEvent; /** @@ -115,38 +114,6 @@ private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundE } } - public static class TestKafkaEvent implements Serializable { - - // Actual content of the event. - private String eventData; - - // Contains Integer value, which is greater than previous message id. - private String eventId; - - TestKafkaEvent(String eventId, String eventData) { - this.eventData = eventData; - this.eventId = eventId; - } - - String getEventId() { - return eventId; - } - - String getEventData() { - return eventData; - } - - @Override - public String toString() { - return eventId + "|" + eventData; - } - - static TestKafkaEvent fromString(String message) { - String[] messageComponents = StringUtils.split(message, "|"); - return new TestKafkaEvent(messageComponents[0], messageComponents[1]); - } - } - public static StreamApplication getInstance( String systemName, List inputTopics, diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java index 4a635cc1aa..feba56e0a4 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java @@ -19,6 +19,7 @@ package org.apache.samza.test.processor; +import java.util.Optional; import java.util.concurrent.CountDownLatch; import org.apache.samza.application.TaskApplication; import org.apache.samza.application.descriptors.TaskApplicationDescriptor; @@ -46,10 +47,32 @@ public class TestTaskApplication implements TaskApplication { private final String outputTopic; private final CountDownLatch shutdownLatch; private final CountDownLatch processedMessageLatch; - private final TaskApplicationCallback processCallback; + private final Optional processCallback; + /** + * A test TaskApplication to use in test harnesses. + * @param systemName test input/output system + * @param inputTopic topic to consume + * @param outputTopic topic to output + * @param processedMessageLatch latch that counts down per message processed + * @param shutdownLatch latch that counts down once during shutdown + */ public TestTaskApplication(String systemName, String inputTopic, String outputTopic, - CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch, TaskApplicationCallback processCallback) { + CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch) { + this(systemName, inputTopic, outputTopic, processedMessageLatch, shutdownLatch, Optional.empty()); + } + + /** + * A test TaskApplication to use in test harnesses. + * @param systemName test input/output system + * @param inputTopic topic to consume + * @param outputTopic topic to output + * @param processedMessageLatch latch that counts down per message processed + * @param shutdownLatch latch that counts down once during shutdown + * @param processCallback optional callback called per message processed. + */ + public TestTaskApplication(String systemName, String inputTopic, String outputTopic, + CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch, Optional processCallback) { this.systemName = systemName; this.inputTopic = inputTopic; this.outputTopic = outputTopic; @@ -65,9 +88,7 @@ public void processAsync(IncomingMessageEnvelope envelope, MessageCollector coll processedMessageLatch.countDown(); // Implementation does not invokes callback.complete to block the RunLoop.process() after it exhausts the // `task.max.concurrency` defined per task. Call callback.complete() in the processCallback if needed. - if (processCallback != null) { - processCallback.onMessage(envelope, callback); - } + processCallback.ifPresent(pcb -> pcb.onMessage(envelope, callback)); } @Override diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java index 01a1279b52..8b7c4042ff 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestZkLocalApplicationRunner.java @@ -24,7 +24,6 @@ import com.google.common.collect.ImmutableSet; import com.google.common.collect.Maps; import com.google.common.collect.Sets; -import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -33,18 +32,13 @@ import java.util.Map; import java.util.Objects; import java.util.Set; -import java.util.TreeMap; import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; import org.I0Itec.zkclient.ZkClient; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.samza.Partition; import org.apache.samza.SamzaException; import org.apache.samza.application.TaskApplication; @@ -71,21 +65,10 @@ import org.apache.samza.runtime.ApplicationRunner; import org.apache.samza.runtime.ApplicationRunners; import org.apache.samza.runtime.LocalApplicationRunner; -import org.apache.samza.startpoint.Startpoint; -import org.apache.samza.startpoint.StartpointManager; -import org.apache.samza.startpoint.StartpointOldest; -import org.apache.samza.startpoint.StartpointSpecific; -import org.apache.samza.startpoint.StartpointTimestamp; -import org.apache.samza.startpoint.StartpointUpcoming; -import org.apache.samza.system.IncomingMessageEnvelope; -import org.apache.samza.system.SystemAdmin; -import org.apache.samza.system.SystemAdmins; -import org.apache.samza.system.SystemStream; import org.apache.samza.system.SystemStreamPartition; -import org.apache.samza.task.TaskCallback; import org.apache.samza.test.StandaloneTestUtils; import org.apache.samza.test.harness.IntegrationTestHarness; -import org.apache.samza.util.CoordinatorStreamUtil; +import org.apache.samza.test.util.TestKafkaEvent; import org.apache.samza.util.NoOpMetricsRegistry; import org.apache.samza.util.ReflectionUtil; import org.apache.samza.zk.ZkJobCoordinatorFactory; @@ -121,7 +104,6 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private static final int NUM_KAFKA_EVENTS = 300; private static final int ZK_CONNECTION_TIMEOUT_MS = 5000; private static final int ZK_SESSION_TIMEOUT_MS = 10000; - private static final int ZK_TEST_PARTITION_COUNT = 5; private static final String TEST_SYSTEM = "TestSystemName"; private static final String TEST_SSP_GROUPER_FACTORY = "org.apache.samza.container.grouper.stream.GroupByPartitionFactory"; private static final String TEST_TASK_GROUPER_FACTORY = "org.apache.samza.container.grouper.task.GroupByContainerIdsFactory"; @@ -130,12 +112,9 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private static final String TASK_SHUTDOWN_MS = "10000"; private static final String JOB_DEBOUNCE_TIME_MS = "10000"; private static final String BARRIER_TIMEOUT_MS = "10000"; - private static final String[] PROCESSOR_IDS = new String[] {"0000000000", "0000000001", "0000000002", "0000000003"}; + private static final String[] PROCESSOR_IDS = new String[] {"0000000000", "0000000001", "0000000002"}; - private String inputKafkaTopic1; - private String inputKafkaTopic2; - private String inputKafkaTopic3; - private String inputKafkaTopic4; + private String inputKafkaTopic; private String outputKafkaTopic; private String inputSinglePartitionKafkaTopic; private String outputSinglePartitionKafkaTopic; @@ -143,7 +122,6 @@ public class TestZkLocalApplicationRunner extends IntegrationTestHarness { private ApplicationConfig applicationConfig1; private ApplicationConfig applicationConfig2; private ApplicationConfig applicationConfig3; - private ApplicationConfig applicationConfig4; private String testStreamAppName; private String testStreamAppId; private MetadataStore zkMetadataStore; @@ -160,10 +138,7 @@ public void setUp() { String uniqueTestId = UUID.randomUUID().toString(); testStreamAppName = String.format("test-app-name-%s", uniqueTestId); testStreamAppId = String.format("test-app-id-%s", uniqueTestId); - inputKafkaTopic1 = String.format("test-input-topic1-%s", uniqueTestId); - inputKafkaTopic2 = String.format("test-input-topic2-%s", uniqueTestId); - inputKafkaTopic3 = String.format("test-input-topic3-%s", uniqueTestId); - inputKafkaTopic4 = String.format("test-input-topic4-%s", uniqueTestId); + inputKafkaTopic = String.format("test-input-topic-%s", uniqueTestId); outputKafkaTopic = String.format("test-output-topic-%s", uniqueTestId); inputSinglePartitionKafkaTopic = String.format("test-input-single-partition-topic-%s", uniqueTestId); outputSinglePartitionKafkaTopic = String.format("test-output-single-partition-topic-%s", uniqueTestId); @@ -178,26 +153,20 @@ public void setUp() { applicationConfig2 = new ApplicationConfig(new MapConfig(configMap)); configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[2]); applicationConfig3 = new ApplicationConfig(new MapConfig(configMap)); - configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[3]); - applicationConfig4 = new ApplicationConfig(new MapConfig(configMap)); ZkClient zkClient = new ZkClient(zkConnect(), ZK_CONNECTION_TIMEOUT_MS, ZK_CONNECTION_TIMEOUT_MS, new ZkStringSerializer()); ZkKeyBuilder zkKeyBuilder = new ZkKeyBuilder(ZkJobCoordinatorFactory.getJobCoordinationZkPath(applicationConfig1)); zkUtils = new ZkUtils(zkKeyBuilder, zkClient, ZK_CONNECTION_TIMEOUT_MS, ZK_SESSION_TIMEOUT_MS, new NoOpMetricsRegistry()); zkUtils.connect(); - ImmutableMap topicToPartitionCount = ImmutableMap.builder() - .put(inputSinglePartitionKafkaTopic, 1) - .put(outputSinglePartitionKafkaTopic, 1) - .put(inputKafkaTopic1, ZK_TEST_PARTITION_COUNT) - .put(inputKafkaTopic2, ZK_TEST_PARTITION_COUNT) - .put(inputKafkaTopic3, ZK_TEST_PARTITION_COUNT) - .put(inputKafkaTopic4, ZK_TEST_PARTITION_COUNT) - .put(outputKafkaTopic, ZK_TEST_PARTITION_COUNT) - .build(); + ImmutableMap topicToPartitionCount = ImmutableMap.of( + inputSinglePartitionKafkaTopic, 1, + outputSinglePartitionKafkaTopic, 1, + inputKafkaTopic, 5, + outputKafkaTopic, 5); List newTopics = - topicToPartitionCount.keySet() + ImmutableList.of(inputKafkaTopic, outputKafkaTopic, inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic) .stream() .map(topic -> new NewTopic(topic, topicToPartitionCount.get(topic), (short) 1)) .collect(Collectors.toList()); @@ -208,8 +177,8 @@ public void setUp() { @Override public void tearDown() { - deleteTopics(ImmutableList.of(inputKafkaTopic1, inputKafkaTopic2, inputKafkaTopic3, inputKafkaTopic4, outputKafkaTopic, - inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic)); + deleteTopics(ImmutableList.of( + inputKafkaTopic, outputKafkaTopic, inputSinglePartitionKafkaTopic, outputSinglePartitionKafkaTopic)); SharedContextFactories.clearAll(); zkUtils.close(); super.tearDown(); @@ -219,36 +188,12 @@ private void publishKafkaEvents(String topic, int startIndex, int endIndex, Stri for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { try { LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); - TestStreamApplication.TestKafkaEvent testKafkaEvent = - new TestStreamApplication.TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)); - producer.send(new ProducerRecord(topic, testKafkaEvent.toString().getBytes())); + producer.send(new ProducerRecord(topic, new TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)).toString().getBytes())); } catch (Exception e) { LOGGER.error("Publishing to kafka topic: {} resulted in exception: {}.", new Object[]{topic, e}); throw new SamzaException(e); } } - producer.flush(); - } - - // Sends each event synchronously and returns map of event index to event metadata of the sent record/event. Adds the specified delay between sends. - private TreeMap publishKafkaEventsWithDelayPerEvent(String topic, int startIndex, int endIndex, String streamProcessorId, Duration delay) { - TreeMap eventsMetadata = new TreeMap<>(); - for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { - try { - LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); - TestStreamApplication.TestKafkaEvent testKafkaEvent = - new TestStreamApplication.TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)); - Thread.sleep(delay.toMillis()); - Future send = producer.send(new ProducerRecord(topic, - testKafkaEvent.toString().getBytes())); - eventsMetadata.put(eventIndex, send.get(5, TimeUnit.SECONDS)); - } catch (Exception e) { - LOGGER.error("Publishing to kafka topic: {} resulted in exception: {}.", new Object[]{topic, e}); - throw new SamzaException(e); - } - } - producer.flush(); - return eventsMetadata; } private Map buildStreamApplicationConfigMap(String appName, String appId, boolean isBatch) { @@ -386,7 +331,7 @@ public void shouldStopNewProcessorsJoiningGroupWhenNumContainersIsGreaterThanNum @Test public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleTaskGrouperFactory() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS * 2, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS * 2, PROCESSOR_IDS[0]); // Configuration, verification variables MapConfig testConfig = new MapConfig(ImmutableMap.of(JobConfig.SSP_GROUPER_FACTORY, @@ -408,7 +353,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT CountDownLatch processedMessagesLatch = new CountDownLatch(NUM_KAFKA_EVENTS * 2); Config testAppConfig2 = new MapConfig(applicationConfig2, testConfig); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch, null, null, testAppConfig2), testAppConfig2); // Callback handler for appRunner1. @@ -432,7 +377,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT // Set up stream app appRunner1. Config testAppConfig1 = new MapConfig(applicationConfig1, testConfig); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, null, streamApplicationCallback, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, null, streamApplicationCallback, kafkaEventsConsumedLatch, testAppConfig1), testAppConfig1); executeRun(appRunner1, testAppConfig1); @@ -477,7 +422,7 @@ public void shouldUpdateJobModelWhenNewProcessorJoiningGroupUsingAllSspToSingleT @Test public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic1, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create stream applications. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(2 * NUM_KAFKA_EVENTS); @@ -486,13 +431,13 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, applicationConfig3), applicationConfig3); executeRun(appRunner1, applicationConfig1); @@ -520,7 +465,7 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { assertEquals(ApplicationStatus.SuccessfulFinish, appRunner1.status()); kafkaEventsConsumedLatch.await(); - publishKafkaEvents(inputKafkaTopic1, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -545,7 +490,7 @@ public void shouldReElectLeaderWhenLeaderDies() throws InterruptedException { @Test public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throws InterruptedException { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create StreamApplications. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); @@ -553,10 +498,10 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Run stream applications. @@ -568,10 +513,10 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw processedMessagesLatch2.await(); // Create a stream app with same processor id as SP2 and run it. It should fail. - publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[2]); + publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[2]); kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, null, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, null, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Fail when the duplicate processor joins. expectedException.expect(SamzaException.class); @@ -586,14 +531,10 @@ public void shouldFailWhenNewProcessorJoinsWithSameIdAsExistingProcessor() throw } } - public void addMessagesProcessed(List messagesProcessed, TestStreamApplication.TestKafkaEvent tke) { - messagesProcessed.add(tke); - } - @Test public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() throws Exception { // Set up kafka topics. - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId); @@ -603,9 +544,8 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[1]); Config applicationConfig2 = new MapConfig(configMap); - List messagesProcessed = new ArrayList<>(); - TestStreamApplication.StreamApplicationCallback streamApplicationCallback = - (TestStreamApplication.TestKafkaEvent tke) -> addMessagesProcessed(messagesProcessed, tke); + List messagesProcessed = new ArrayList<>(); + TestStreamApplication.StreamApplicationCallback streamApplicationCallback = messagesProcessed::add; // Create StreamApplication from configuration. CountDownLatch kafkaEventsConsumedLatch = new CountDownLatch(NUM_KAFKA_EVENTS); @@ -613,10 +553,10 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, streamApplicationCallback, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); // Run stream application. @@ -634,7 +574,7 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t appRunner1.waitForFinish(); int lastProcessedMessageId = -1; - for (TestStreamApplication.TestKafkaEvent message : messagesProcessed) { + for (TestKafkaEvent message : messagesProcessed) { lastProcessedMessageId = Math.max(lastProcessedMessageId, Integer.parseInt(message.getEventData())); } messagesProcessed.clear(); @@ -642,9 +582,9 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t assertEquals(ApplicationStatus.SuccessfulFinish, appRunner1.status()); processedMessagesLatch1 = new CountDownLatch(1); - publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); executeRun(appRunner3, applicationConfig1); @@ -667,7 +607,7 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t @Test public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContainerShutdownTime() throws Exception { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId); configMap.put(TaskConfig.TASK_SHUTDOWN_MS, "0"); @@ -684,10 +624,10 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, kafkaEventsConsumedLatch, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -709,11 +649,11 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, kafkaEventsConsumedLatch, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); - publishKafkaEvents(inputKafkaTopic1, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, NUM_KAFKA_EVENTS, 2 * NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); processedMessagesLatch3.await(); appRunner1.waitForFinish(); @@ -739,14 +679,14 @@ public void testShouldStopStreamApplicationWhenShutdownTimeOutIsLessThanContaine */ @Test public void testShouldGenerateJobModelOnPartitionCountChange() throws Exception { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create StreamApplication from configuration. CountDownLatch kafkaEventsConsumedLatch1 = new CountDownLatch(NUM_KAFKA_EVENTS); CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch1, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, kafkaEventsConsumedLatch1, applicationConfig1), applicationConfig1); executeRun(appRunner1, applicationConfig1); @@ -760,7 +700,7 @@ public void testShouldGenerateJobModelOnPartitionCountChange() throws Exception Assert.assertEquals(5, ssps.size()); // Increase the partition count of input kafka topic to 100. - increasePartitionsTo(inputKafkaTopic1, 100); + increasePartitionsTo(inputKafkaTopic, 100); long jobModelWaitTimeInMillis = 10; while (Objects.equals(zkUtils.getJobModelVersion(), jobModelVersion)) { @@ -1002,14 +942,14 @@ public void testStatefulSamzaApplicationShouldRedistributeInputPartitionsToCorre @Test public void testApplicationShutdownShouldBeIndependentOfPerMessageProcessingTime() throws Exception { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); // Create a TaskApplication with only one task per container. // The task does not invokes taskCallback.complete for any of the dispatched message. CountDownLatch shutdownLatch = new CountDownLatch(1); CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); - TaskApplication taskApplication = new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatch1, shutdownLatch, null); + TaskApplication taskApplication = new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic, outputKafkaTopic, processedMessagesLatch1, shutdownLatch); MapConfig taskApplicationConfig = new MapConfig(ImmutableList.of(applicationConfig1, ImmutableMap.of(TaskConfig.MAX_CONCURRENCY, "1", JobConfig.SSP_GROUPER_FACTORY, "org.apache.samza.container.grouper.stream.AllSspToSingleTaskGrouperFactory"))); ApplicationRunner appRunner = ApplicationRunners.getApplicationRunner(taskApplication, taskApplicationConfig); @@ -1039,7 +979,7 @@ public void testApplicationShutdownShouldBeIndependentOfPerMessageProcessingTime */ @Test public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1054,10 +994,10 @@ public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1091,7 +1031,7 @@ public void testAgreeingOnSameRunIdForBatch() throws InterruptedException { */ @Test public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1107,10 +1047,10 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1130,7 +1070,7 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1161,7 +1101,7 @@ public void testNewProcessorGetsSameRunIdForBatch() throws InterruptedException */ @Test public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1177,10 +1117,10 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner1, applicationConfig1); @@ -1208,7 +1148,7 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1235,7 +1175,7 @@ public void testAllProcesssorDieNewProcessorGetsNewRunIdForBatch() throws Interr */ @Test public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedException { - publishKafkaEvents(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); + publishKafkaEvents(inputKafkaTopic, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0]); Map configMap = buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId, true); @@ -1248,7 +1188,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep CountDownLatch processedMessagesLatch1 = new CountDownLatch(1); ApplicationRunner appRunner1 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch1, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch1, null, null, applicationConfig1), applicationConfig1); executeRun(appRunner1, applicationConfig1); @@ -1260,7 +1200,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep // bring up second processor CountDownLatch processedMessagesLatch2 = new CountDownLatch(1); ApplicationRunner appRunner2 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch2, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch2, null, null, applicationConfig2), applicationConfig2); executeRun(appRunner2, applicationConfig2); @@ -1277,7 +1217,7 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep //Bring up a new processsor CountDownLatch processedMessagesLatch3 = new CountDownLatch(1); ApplicationRunner appRunner3 = ApplicationRunners.getApplicationRunner(TestStreamApplication.getInstance( - TEST_SYSTEM, ImmutableList.of(inputKafkaTopic1), outputKafkaTopic, processedMessagesLatch3, null, null, + TEST_SYSTEM, ImmutableList.of(inputKafkaTopic), outputKafkaTopic, processedMessagesLatch3, null, null, applicationConfig3), applicationConfig3); executeRun(appRunner3, applicationConfig3); processedMessagesLatch3.await(); @@ -1294,159 +1234,6 @@ public void testFirstProcessorDiesButSameRunIdForBatch() throws InterruptedExcep appRunner3.waitForFinish(); } - private void writeStartpoints(StartpointManager startpointManager, String streamName, Integer partitionCount, Startpoint startpoint) { - for (int p = 0; p < partitionCount; p++) { - startpointManager.writeStartpoint(new SystemStreamPartition(TEST_SYSTEM, streamName, new Partition(p)), startpoint); - } - } - - @Test - public void testStartpoints() throws InterruptedException { - TreeMap sentEvents1 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0], Duration.ofMillis(2)); - TreeMap sentEvents2 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic2, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[1], Duration.ofMillis(2)); - TreeMap sentEvents3 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic3, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[2], Duration.ofMillis(2)); - TreeMap sentEvents4 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic4, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[3], Duration.ofMillis(2)); - - ConcurrentHashMap recvEventsInputStartpointSpecific = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointTimestamp = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointOldest = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointUpcoming = new ConcurrentHashMap<>(); - - CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); - coordinatorStreamStore.init(); - StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); - startpointManager.start(); - - StartpointSpecific startpointSpecific = new StartpointSpecific(String.valueOf(sentEvents1.get(100).offset())); - writeStartpoints(startpointManager, inputKafkaTopic1, ZK_TEST_PARTITION_COUNT, startpointSpecific); - - StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); - writeStartpoints(startpointManager, inputKafkaTopic2, ZK_TEST_PARTITION_COUNT, startpointTimestamp); - - StartpointOldest startpointOldest = new StartpointOldest(); - writeStartpoints(startpointManager, inputKafkaTopic3, ZK_TEST_PARTITION_COUNT, startpointOldest); - - StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); - writeStartpoints(startpointManager, inputKafkaTopic4, ZK_TEST_PARTITION_COUNT, startpointUpcoming); - - startpointManager.stop(); - coordinatorStreamStore.close(); - - TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { - try { - String streamName = ime.getSystemStreamPartition().getStream(); - TestStreamApplication.TestKafkaEvent testKafkaEvent = - TestStreamApplication.TestKafkaEvent.fromString((String) ime.getMessage()); - String eventIndex = testKafkaEvent.getEventData(); - if (inputKafkaTopic1.equals(streamName)) { - recvEventsInputStartpointSpecific.put(eventIndex, ime); - } else if (inputKafkaTopic2.equals(streamName)) { - recvEventsInputStartpointTimestamp.put(eventIndex, ime); - } else if (inputKafkaTopic3.equals(streamName)) { - recvEventsInputStartpointOldest.put(eventIndex, ime); - } else if (inputKafkaTopic4.equals(streamName)) { - recvEventsInputStartpointUpcoming.put(eventIndex, ime); - } else { - throw new RuntimeException("Unexpected input stream: " + streamName); - } - callback.complete(); - } catch (Exception ex) { - callback.failure(ex); - } - }; - - // Create StreamApplication from configuration. - CountDownLatch processedMessagesLatchStartpointSpecific = new CountDownLatch(100); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Fetch all since consuming from oldest - CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Expecting none, so just attempt a small number of fetches. - - CountDownLatch shutdownLatchStartpointSpecific = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointOldest = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointUpcoming = new CountDownLatch(1); - - TestTaskApplication testTaskApplicationStartpointSpecific = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatchStartpointSpecific, - shutdownLatchStartpointSpecific, processedCallback); - TestTaskApplication testTaskApplicationStartpointTimestamp = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic2, outputKafkaTopic, processedMessagesLatchStartpointTimestamp, - shutdownLatchStartpointTimestamp, processedCallback); - TestTaskApplication testTaskApplicationStartpointOldest = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic3, outputKafkaTopic, processedMessagesLatchStartpointOldest, - shutdownLatchStartpointOldest, processedCallback); - TestTaskApplication testTaskApplicationStartpointUpcoming = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic4, outputKafkaTopic, processedMessagesLatchStartpointUpcoming, - shutdownLatchStartpointUpcoming, processedCallback); - - // Startpoint specific - ApplicationRunner appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointSpecific, applicationConfig1); - executeRun(appRunner, applicationConfig1); - - assertTrue(processedMessagesLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); - - appRunner.kill(); - appRunner.waitForFinish(); - - assertTrue(shutdownLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); - - Integer startpointSpecificOffset = Integer.valueOf(startpointSpecific.getSpecificOffset()); - for (IncomingMessageEnvelope ime : recvEventsInputStartpointSpecific.values()) { - Integer eventOffset = Integer.valueOf(ime.getOffset()); - String assertMsg = String.format("Expecting message offset: %d >= Startpoint specific offset: %d", - eventOffset, startpointSpecificOffset); - assertTrue(assertMsg, eventOffset >= startpointSpecificOffset); - } - - // Startpoint timestamp - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointTimestamp, applicationConfig2); - - executeRun(appRunner, applicationConfig2); - assertTrue(processedMessagesLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); - - appRunner.kill(); - appRunner.waitForFinish(); - - assertTrue(shutdownLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); - - for (IncomingMessageEnvelope ime : recvEventsInputStartpointTimestamp.values()) { - Integer eventOffset = Integer.valueOf(ime.getOffset()); - assertNotEquals(0, ime.getEventTime()); // sanity check - String assertMsg = String.format("Expecting message timestamp: %d >= Startpoint timestamp: %d", - ime.getEventTime(), startpointTimestamp.getTimestampOffset()); - assertTrue(assertMsg, ime.getEventTime() >= startpointTimestamp.getTimestampOffset()); - } - - // Startpoint oldest - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointOldest, applicationConfig3); - executeRun(appRunner, applicationConfig3); - - assertTrue(processedMessagesLatchStartpointOldest.await(1, TimeUnit.MINUTES)); - - appRunner.kill(); - appRunner.waitForFinish(); - - assertTrue(shutdownLatchStartpointOldest.await(1, TimeUnit.MINUTES)); - - assertEquals("Expecting to have processed all the events", NUM_KAFKA_EVENTS, recvEventsInputStartpointOldest.size()); - - // Startpoint upcoming - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointUpcoming, applicationConfig4); - executeRun(appRunner, applicationConfig4); - - assertFalse("Expecting to timeout and not process any old messages.", processedMessagesLatchStartpointUpcoming.await(15, TimeUnit.SECONDS)); - assertEquals("Expecting not to process any old messages.", 0, recvEventsInputStartpointUpcoming.size()); - - appRunner.kill(); - appRunner.waitForFinish(); - - assertTrue(shutdownLatchStartpointUpcoming.await(1, TimeUnit.MINUTES)); - } - /** * Computes the task to partition assignment of the {@param JobModel}. * @param jobModel the jobModel to compute task to partition assignment for. @@ -1475,14 +1262,4 @@ private static List getSystemStreamPartitions(JobModel jo }); return ssps; } - - private static CoordinatorStreamStore createCoordinatorStreamStore(Config applicationConfig) { - SystemStream coordinatorSystemStream = CoordinatorStreamUtil.getCoordinatorSystemStream(applicationConfig); - SystemAdmins systemAdmins = new SystemAdmins(applicationConfig); - SystemAdmin coordinatorSystemAdmin = systemAdmins.getSystemAdmin(coordinatorSystemStream.getSystem()); - coordinatorSystemAdmin.start(); - CoordinatorStreamUtil.createCoordinatorStream(coordinatorSystemStream, coordinatorSystemAdmin); - coordinatorSystemAdmin.stop(); - return new CoordinatorStreamStore(applicationConfig, new NoOpMetricsRegistry()); - } } diff --git a/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java new file mode 100644 index 0000000000..8b2a05b900 --- /dev/null +++ b/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java @@ -0,0 +1,383 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.samza.test.startpoint; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.Maps; +import java.time.Duration; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.TreeMap; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.samza.Partition; +import org.apache.samza.SamzaException; +import org.apache.samza.config.ApplicationConfig; +import org.apache.samza.config.Config; +import org.apache.samza.config.JobConfig; +import org.apache.samza.config.JobCoordinatorConfig; +import org.apache.samza.config.MapConfig; +import org.apache.samza.config.TaskConfig; +import org.apache.samza.config.ZkConfig; +import org.apache.samza.coordinator.metadatastore.CoordinatorStreamStore; +import org.apache.samza.runtime.ApplicationRunner; +import org.apache.samza.runtime.ApplicationRunners; +import org.apache.samza.startpoint.Startpoint; +import org.apache.samza.startpoint.StartpointManager; +import org.apache.samza.startpoint.StartpointOldest; +import org.apache.samza.startpoint.StartpointSpecific; +import org.apache.samza.startpoint.StartpointTimestamp; +import org.apache.samza.startpoint.StartpointUpcoming; +import org.apache.samza.system.IncomingMessageEnvelope; +import org.apache.samza.system.SystemAdmin; +import org.apache.samza.system.SystemAdmins; +import org.apache.samza.system.SystemStream; +import org.apache.samza.system.SystemStreamPartition; +import org.apache.samza.task.TaskCallback; +import org.apache.samza.test.StandaloneTestUtils; +import org.apache.samza.test.harness.IntegrationTestHarness; +import org.apache.samza.test.processor.SharedContextFactories; +import org.apache.samza.test.processor.TestTaskApplication; +import org.apache.samza.test.util.TestKafkaEvent; +import org.apache.samza.util.CoordinatorStreamUtil; +import org.apache.samza.util.NoOpMetricsRegistry; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.ExpectedException; +import org.junit.rules.Timeout; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertTrue; + + +public class StartpointTestHarness extends IntegrationTestHarness { + + private static final Logger LOGGER = LoggerFactory.getLogger(StartpointTestHarness.class); + + private static final int NUM_KAFKA_EVENTS = 300; + private static final int ZK_TEST_PARTITION_COUNT = 5; + private static final String TEST_SYSTEM = "TestSystemName"; + private static final String TEST_SSP_GROUPER_FACTORY = "org.apache.samza.container.grouper.stream.GroupByPartitionFactory"; + private static final String TEST_TASK_GROUPER_FACTORY = "org.apache.samza.container.grouper.task.GroupByContainerIdsFactory"; + private static final String TEST_JOB_COORDINATOR_FACTORY = "org.apache.samza.zk.ZkJobCoordinatorFactory"; + private static final String TEST_SYSTEM_FACTORY = "org.apache.samza.system.kafka.KafkaSystemFactory"; + private static final String TASK_SHUTDOWN_MS = "10000"; + private static final String JOB_DEBOUNCE_TIME_MS = "10000"; + private static final String BARRIER_TIMEOUT_MS = "10000"; + private static final String[] PROCESSOR_IDS = new String[] {"0000000000", "0000000001", "0000000002", "0000000003"}; + + private String inputKafkaTopic1; + private String inputKafkaTopic2; + private String inputKafkaTopic3; + private String inputKafkaTopic4; + private String outputKafkaTopic; + private ApplicationConfig applicationConfig1; + private ApplicationConfig applicationConfig2; + private ApplicationConfig applicationConfig3; + private ApplicationConfig applicationConfig4; + private String testStreamAppName; + private String testStreamAppId; + + @Rule + public Timeout testTimeOutInMillis = new Timeout(150000, TimeUnit.MILLISECONDS); + + @Rule + public final ExpectedException expectedException = ExpectedException.none(); + + @Override + public void setUp() { + super.setUp(); + String uniqueTestId = UUID.randomUUID().toString(); + testStreamAppName = String.format("test-app-name-%s", uniqueTestId); + testStreamAppId = String.format("test-app-id-%s", uniqueTestId); + inputKafkaTopic1 = String.format("test-input-topic1-%s", uniqueTestId); + inputKafkaTopic2 = String.format("test-input-topic2-%s", uniqueTestId); + inputKafkaTopic3 = String.format("test-input-topic3-%s", uniqueTestId); + inputKafkaTopic4 = String.format("test-input-topic4-%s", uniqueTestId); + outputKafkaTopic = String.format("test-output-topic-%s", uniqueTestId); + + // Set up stream application config map with the given testStreamAppName, testStreamAppId and test kafka system + // TODO: processorId should typically come up from a processorID generator as processor.id will be deprecated in 0.14.0+ + Map configMap = + buildStreamApplicationConfigMap(testStreamAppName, testStreamAppId); + configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[0]); + applicationConfig1 = new ApplicationConfig(new MapConfig(configMap)); + configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[1]); + applicationConfig2 = new ApplicationConfig(new MapConfig(configMap)); + configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[2]); + applicationConfig3 = new ApplicationConfig(new MapConfig(configMap)); + configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[3]); + applicationConfig4 = new ApplicationConfig(new MapConfig(configMap)); + + ImmutableMap topicToPartitionCount = ImmutableMap.builder() + .put(inputKafkaTopic1, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic2, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic3, ZK_TEST_PARTITION_COUNT) + .put(inputKafkaTopic4, ZK_TEST_PARTITION_COUNT) + .put(outputKafkaTopic, ZK_TEST_PARTITION_COUNT) + .build(); + + List newTopics = + topicToPartitionCount.keySet() + .stream() + .map(topic -> new NewTopic(topic, topicToPartitionCount.get(topic), (short) 1)) + .collect(Collectors.toList()); + + assertTrue("Encountered errors during test setup. Failed to create topics.", createTopics(newTopics)); + } + + @Override + public void tearDown() { + deleteTopics( + ImmutableList.of(inputKafkaTopic1, inputKafkaTopic2, inputKafkaTopic3, inputKafkaTopic4, outputKafkaTopic)); + SharedContextFactories.clearAll(); + super.tearDown(); + } + + @Test + public void testStartpoints() throws InterruptedException { + TreeMap sentEvents1 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0], Duration.ofMillis(2)); + TreeMap sentEvents2 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic2, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[1], Duration.ofMillis(2)); + TreeMap sentEvents3 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic3, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[2], Duration.ofMillis(2)); + TreeMap sentEvents4 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic4, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[3], Duration.ofMillis(2)); + + ConcurrentHashMap recvEventsInputStartpointSpecific = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointTimestamp = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointOldest = new ConcurrentHashMap<>(); + ConcurrentHashMap recvEventsInputStartpointUpcoming = new ConcurrentHashMap<>(); + + CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); + coordinatorStreamStore.init(); + StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); + startpointManager.start(); + + StartpointSpecific startpointSpecific = new StartpointSpecific(String.valueOf(sentEvents1.get(100).offset())); + writeStartpoints(startpointManager, inputKafkaTopic1, ZK_TEST_PARTITION_COUNT, startpointSpecific); + + StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); + writeStartpoints(startpointManager, inputKafkaTopic2, ZK_TEST_PARTITION_COUNT, startpointTimestamp); + + StartpointOldest startpointOldest = new StartpointOldest(); + writeStartpoints(startpointManager, inputKafkaTopic3, ZK_TEST_PARTITION_COUNT, startpointOldest); + + StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); + writeStartpoints(startpointManager, inputKafkaTopic4, ZK_TEST_PARTITION_COUNT, startpointUpcoming); + + startpointManager.stop(); + coordinatorStreamStore.close(); + + TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + try { + String streamName = ime.getSystemStreamPartition().getStream(); + TestKafkaEvent testKafkaEvent = + TestKafkaEvent.fromString((String) ime.getMessage()); + String eventIndex = testKafkaEvent.getEventData(); + if (inputKafkaTopic1.equals(streamName)) { + recvEventsInputStartpointSpecific.put(eventIndex, ime); + } else if (inputKafkaTopic2.equals(streamName)) { + recvEventsInputStartpointTimestamp.put(eventIndex, ime); + } else if (inputKafkaTopic3.equals(streamName)) { + recvEventsInputStartpointOldest.put(eventIndex, ime); + } else if (inputKafkaTopic4.equals(streamName)) { + recvEventsInputStartpointUpcoming.put(eventIndex, ime); + } else { + throw new RuntimeException("Unexpected input stream: " + streamName); + } + callback.complete(); + } catch (Exception ex) { + callback.failure(ex); + } + }; + + // Create StreamApplication from configuration. + CountDownLatch processedMessagesLatchStartpointSpecific = new CountDownLatch(100); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages + CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Fetch all since consuming from oldest + CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Expecting none, so just attempt a small number of fetches. + + CountDownLatch shutdownLatchStartpointSpecific = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointOldest = new CountDownLatch(1); + CountDownLatch shutdownLatchStartpointUpcoming = new CountDownLatch(1); + + TestTaskApplication testTaskApplicationStartpointSpecific = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatchStartpointSpecific, + shutdownLatchStartpointSpecific, Optional.of(processedCallback)); + TestTaskApplication testTaskApplicationStartpointTimestamp = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic2, outputKafkaTopic, processedMessagesLatchStartpointTimestamp, + shutdownLatchStartpointTimestamp, Optional.of(processedCallback)); + TestTaskApplication testTaskApplicationStartpointOldest = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic3, outputKafkaTopic, processedMessagesLatchStartpointOldest, + shutdownLatchStartpointOldest, Optional.of(processedCallback)); + TestTaskApplication testTaskApplicationStartpointUpcoming = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic4, outputKafkaTopic, processedMessagesLatchStartpointUpcoming, + shutdownLatchStartpointUpcoming, Optional.of(processedCallback)); + + // Startpoint specific + ApplicationRunner + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointSpecific, applicationConfig1); + executeRun(appRunner, applicationConfig1); + + assertTrue(processedMessagesLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); + + Integer startpointSpecificOffset = Integer.valueOf(startpointSpecific.getSpecificOffset()); + for (IncomingMessageEnvelope ime : recvEventsInputStartpointSpecific.values()) { + Integer eventOffset = Integer.valueOf(ime.getOffset()); + String assertMsg = String.format("Expecting message offset: %d >= Startpoint specific offset: %d", + eventOffset, startpointSpecificOffset); + assertTrue(assertMsg, eventOffset >= startpointSpecificOffset); + } + + // Startpoint timestamp + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointTimestamp, applicationConfig2); + + executeRun(appRunner, applicationConfig2); + assertTrue(processedMessagesLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); + + for (IncomingMessageEnvelope ime : recvEventsInputStartpointTimestamp.values()) { + Integer eventOffset = Integer.valueOf(ime.getOffset()); + assertNotEquals(0, ime.getEventTime()); // sanity check + String assertMsg = String.format("Expecting message timestamp: %d >= Startpoint timestamp: %d", + ime.getEventTime(), startpointTimestamp.getTimestampOffset()); + assertTrue(assertMsg, ime.getEventTime() >= startpointTimestamp.getTimestampOffset()); + } + + // Startpoint oldest + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointOldest, applicationConfig3); + executeRun(appRunner, applicationConfig3); + + assertTrue(processedMessagesLatchStartpointOldest.await(1, TimeUnit.MINUTES)); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointOldest.await(1, TimeUnit.MINUTES)); + + assertEquals("Expecting to have processed all the events", NUM_KAFKA_EVENTS, recvEventsInputStartpointOldest.size()); + + // Startpoint upcoming + appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointUpcoming, applicationConfig4); + executeRun(appRunner, applicationConfig4); + + assertFalse("Expecting to timeout and not process any old messages.", processedMessagesLatchStartpointUpcoming.await(15, TimeUnit.SECONDS)); + assertEquals("Expecting not to process any old messages.", 0, recvEventsInputStartpointUpcoming.size()); + + appRunner.kill(); + appRunner.waitForFinish(); + + assertTrue(shutdownLatchStartpointUpcoming.await(1, TimeUnit.MINUTES)); + } + + // Sends each event synchronously and returns map of event index to event metadata of the sent record/event. Adds the specified delay between sends. + private TreeMap publishKafkaEventsWithDelayPerEvent(String topic, int startIndex, int endIndex, String streamProcessorId, Duration delay) { + TreeMap eventsMetadata = new TreeMap<>(); + for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { + try { + LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); + TestKafkaEvent testKafkaEvent = + new TestKafkaEvent(streamProcessorId, String.valueOf(eventIndex)); + Thread.sleep(delay.toMillis()); + Future send = producer.send(new ProducerRecord(topic, + testKafkaEvent.toString().getBytes())); + eventsMetadata.put(eventIndex, send.get(5, TimeUnit.SECONDS)); + } catch (Exception e) { + LOGGER.error("Publishing to kafka topic: {} resulted in exception: {}.", new Object[]{topic, e}); + throw new SamzaException(e); + } + } + producer.flush(); + return eventsMetadata; + } + + private void writeStartpoints(StartpointManager startpointManager, String streamName, Integer partitionCount, Startpoint startpoint) { + for (int p = 0; p < partitionCount; p++) { + startpointManager.writeStartpoint(new SystemStreamPartition(TEST_SYSTEM, streamName, new Partition(p)), startpoint); + } + } + + private static CoordinatorStreamStore createCoordinatorStreamStore(Config applicationConfig) { + SystemStream coordinatorSystemStream = CoordinatorStreamUtil.getCoordinatorSystemStream(applicationConfig); + SystemAdmins systemAdmins = new SystemAdmins(applicationConfig); + SystemAdmin coordinatorSystemAdmin = systemAdmins.getSystemAdmin(coordinatorSystemStream.getSystem()); + coordinatorSystemAdmin.start(); + CoordinatorStreamUtil.createCoordinatorStream(coordinatorSystemStream, coordinatorSystemAdmin); + coordinatorSystemAdmin.stop(); + return new CoordinatorStreamStore(applicationConfig, new NoOpMetricsRegistry()); + } + + private Map buildStreamApplicationConfigMap(String appName, String appId) { + String coordinatorSystemName = "coordinatorSystem"; + Map config = new HashMap<>(); + config.put(ZkConfig.ZK_CONSENSUS_TIMEOUT_MS, BARRIER_TIMEOUT_MS); + config.put(JobConfig.JOB_DEFAULT_SYSTEM, TEST_SYSTEM); + config.put(TaskConfig.IGNORED_EXCEPTIONS, "*"); + config.put(ZkConfig.ZK_CONNECT, zkConnect()); + config.put(JobConfig.SSP_GROUPER_FACTORY, TEST_SSP_GROUPER_FACTORY); + config.put(TaskConfig.GROUPER_FACTORY, TEST_TASK_GROUPER_FACTORY); + config.put(JobCoordinatorConfig.JOB_COORDINATOR_FACTORY, TEST_JOB_COORDINATOR_FACTORY); + config.put(ApplicationConfig.APP_NAME, appName); + config.put(ApplicationConfig.APP_ID, appId); + config.put("app.runner.class", "org.apache.samza.runtime.LocalApplicationRunner"); + config.put(String.format("systems.%s.samza.factory", TEST_SYSTEM), TEST_SYSTEM_FACTORY); + config.put(JobConfig.JOB_NAME, appName); + config.put(JobConfig.JOB_ID, appId); + config.put(TaskConfig.TASK_SHUTDOWN_MS, TASK_SHUTDOWN_MS); + config.put(TaskConfig.DROP_PRODUCER_ERRORS, "true"); + config.put(JobConfig.JOB_DEBOUNCE_TIME_MS, JOB_DEBOUNCE_TIME_MS); + config.put(JobConfig.MONITOR_PARTITION_CHANGE_FREQUENCY_MS, "1000"); + config.put("job.coordinator.system", coordinatorSystemName); + config.put("job.coordinator.replication.factor", "1"); + + Map samzaContainerConfig = ImmutableMap.builder().putAll(config).build(); + Map applicationConfig = Maps.newHashMap(samzaContainerConfig); + applicationConfig.putAll( + StandaloneTestUtils.getKafkaSystemConfigs(coordinatorSystemName, bootstrapServers(), zkConnect(), null, StandaloneTestUtils.SerdeAlias.STRING, true)); + applicationConfig.putAll(StandaloneTestUtils.getKafkaSystemConfigs(TEST_SYSTEM, bootstrapServers(), zkConnect(), null, StandaloneTestUtils.SerdeAlias.STRING, true)); + return applicationConfig; + } +} diff --git a/samza-test/src/test/java/org/apache/samza/test/util/TestKafkaEvent.java b/samza-test/src/test/java/org/apache/samza/test/util/TestKafkaEvent.java new file mode 100644 index 0000000000..575e526841 --- /dev/null +++ b/samza-test/src/test/java/org/apache/samza/test/util/TestKafkaEvent.java @@ -0,0 +1,58 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.samza.test.util; + +import java.io.Serializable; +import org.apache.commons.lang3.StringUtils; + + +/** + * A Kafka event to use in testing only. + */ +public class TestKafkaEvent implements Serializable { + // Actual content of the event. + private String eventData; + + // Contains Integer value, which is greater than previous message id. + private String eventId; + + public TestKafkaEvent(String eventId, String eventData) { + this.eventData = eventData; + this.eventId = eventId; + } + + public String getEventId() { + return eventId; + } + + public String getEventData() { + return eventData; + } + + @Override + public String toString() { + return eventId + "|" + eventData; + } + + public static TestKafkaEvent fromString(String message) { + String[] messageComponents = StringUtils.split(message, "|"); + return new TestKafkaEvent(messageComponents[0], messageComponents[1]); + } +} From 6e748d8de3de343e2094a5e30e5d736131750cf0 Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Fri, 27 Sep 2019 16:05:39 -0700 Subject: [PATCH 4/6] Address @PawasChhokra review comments. --- .../startpoint/StartpointTestHarness.java | 205 +++++++++++++----- 1 file changed, 145 insertions(+), 60 deletions(-) diff --git a/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java index 8b2a05b900..c064b7bdfd 100644 --- a/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java +++ b/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java @@ -27,7 +27,6 @@ import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.TreeMap; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; @@ -166,20 +165,11 @@ public void tearDown() { } @Test - public void testStartpoints() throws InterruptedException { - TreeMap sentEvents1 = + public void testStartpointSpecific() throws InterruptedException { + Map sentEvents1 = publishKafkaEventsWithDelayPerEvent(inputKafkaTopic1, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[0], Duration.ofMillis(2)); - TreeMap sentEvents2 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic2, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[1], Duration.ofMillis(2)); - TreeMap sentEvents3 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic3, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[2], Duration.ofMillis(2)); - TreeMap sentEvents4 = - publishKafkaEventsWithDelayPerEvent(inputKafkaTopic4, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[3], Duration.ofMillis(2)); ConcurrentHashMap recvEventsInputStartpointSpecific = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointTimestamp = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointOldest = new ConcurrentHashMap<>(); - ConcurrentHashMap recvEventsInputStartpointUpcoming = new ConcurrentHashMap<>(); CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); coordinatorStreamStore.init(); @@ -189,32 +179,16 @@ public void testStartpoints() throws InterruptedException { StartpointSpecific startpointSpecific = new StartpointSpecific(String.valueOf(sentEvents1.get(100).offset())); writeStartpoints(startpointManager, inputKafkaTopic1, ZK_TEST_PARTITION_COUNT, startpointSpecific); - StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); - writeStartpoints(startpointManager, inputKafkaTopic2, ZK_TEST_PARTITION_COUNT, startpointTimestamp); - - StartpointOldest startpointOldest = new StartpointOldest(); - writeStartpoints(startpointManager, inputKafkaTopic3, ZK_TEST_PARTITION_COUNT, startpointOldest); - - StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); - writeStartpoints(startpointManager, inputKafkaTopic4, ZK_TEST_PARTITION_COUNT, startpointUpcoming); - startpointManager.stop(); coordinatorStreamStore.close(); TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { try { String streamName = ime.getSystemStreamPartition().getStream(); - TestKafkaEvent testKafkaEvent = - TestKafkaEvent.fromString((String) ime.getMessage()); + TestKafkaEvent testKafkaEvent = TestKafkaEvent.fromString((String) ime.getMessage()); String eventIndex = testKafkaEvent.getEventData(); if (inputKafkaTopic1.equals(streamName)) { recvEventsInputStartpointSpecific.put(eventIndex, ime); - } else if (inputKafkaTopic2.equals(streamName)) { - recvEventsInputStartpointTimestamp.put(eventIndex, ime); - } else if (inputKafkaTopic3.equals(streamName)) { - recvEventsInputStartpointOldest.put(eventIndex, ime); - } else if (inputKafkaTopic4.equals(streamName)) { - recvEventsInputStartpointUpcoming.put(eventIndex, ime); } else { throw new RuntimeException("Unexpected input stream: " + streamName); } @@ -224,33 +198,15 @@ public void testStartpoints() throws InterruptedException { } }; - // Create StreamApplication from configuration. CountDownLatch processedMessagesLatchStartpointSpecific = new CountDownLatch(100); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages - CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Fetch all since consuming from oldest - CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Expecting none, so just attempt a small number of fetches. - CountDownLatch shutdownLatchStartpointSpecific = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointOldest = new CountDownLatch(1); - CountDownLatch shutdownLatchStartpointUpcoming = new CountDownLatch(1); TestTaskApplication testTaskApplicationStartpointSpecific = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, processedMessagesLatchStartpointSpecific, - shutdownLatchStartpointSpecific, Optional.of(processedCallback)); - TestTaskApplication testTaskApplicationStartpointTimestamp = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic2, outputKafkaTopic, processedMessagesLatchStartpointTimestamp, - shutdownLatchStartpointTimestamp, Optional.of(processedCallback)); - TestTaskApplication testTaskApplicationStartpointOldest = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic3, outputKafkaTopic, processedMessagesLatchStartpointOldest, - shutdownLatchStartpointOldest, Optional.of(processedCallback)); - TestTaskApplication testTaskApplicationStartpointUpcoming = - new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic4, outputKafkaTopic, processedMessagesLatchStartpointUpcoming, - shutdownLatchStartpointUpcoming, Optional.of(processedCallback)); + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic1, outputKafkaTopic, + processedMessagesLatchStartpointSpecific, shutdownLatchStartpointSpecific, Optional.of(processedCallback)); - // Startpoint specific - ApplicationRunner - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointSpecific, applicationConfig1); + ApplicationRunner appRunner = + ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointSpecific, applicationConfig1); executeRun(appRunner, applicationConfig1); assertTrue(processedMessagesLatchStartpointSpecific.await(1, TimeUnit.MINUTES)); @@ -263,13 +219,56 @@ public void testStartpoints() throws InterruptedException { Integer startpointSpecificOffset = Integer.valueOf(startpointSpecific.getSpecificOffset()); for (IncomingMessageEnvelope ime : recvEventsInputStartpointSpecific.values()) { Integer eventOffset = Integer.valueOf(ime.getOffset()); - String assertMsg = String.format("Expecting message offset: %d >= Startpoint specific offset: %d", - eventOffset, startpointSpecificOffset); + String assertMsg = String.format("Expecting message offset: %d >= Startpoint specific offset: %d", eventOffset, + startpointSpecificOffset); assertTrue(assertMsg, eventOffset >= startpointSpecificOffset); } + } - // Startpoint timestamp - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointTimestamp, applicationConfig2); + @Test + public void testStartpointTimestamp() throws InterruptedException { + Map sentEvents2 = + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic2, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[1], Duration.ofMillis(2)); + + ConcurrentHashMap recvEventsInputStartpointTimestamp = new ConcurrentHashMap<>(); + + CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); + coordinatorStreamStore.init(); + StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); + startpointManager.start(); + + StartpointTimestamp startpointTimestamp = new StartpointTimestamp(sentEvents2.get(150).timestamp()); + writeStartpoints(startpointManager, inputKafkaTopic2, ZK_TEST_PARTITION_COUNT, startpointTimestamp); + + startpointManager.stop(); + coordinatorStreamStore.close(); + + TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + try { + String streamName = ime.getSystemStreamPartition().getStream(); + TestKafkaEvent testKafkaEvent = + TestKafkaEvent.fromString((String) ime.getMessage()); + String eventIndex = testKafkaEvent.getEventData(); + if (inputKafkaTopic2.equals(streamName)) { + recvEventsInputStartpointTimestamp.put(eventIndex, ime); + } else { + throw new RuntimeException("Unexpected input stream: " + streamName); + } + callback.complete(); + } catch (Exception ex) { + callback.failure(ex); + } + }; + + CountDownLatch processedMessagesLatchStartpointTimestamp = new CountDownLatch(100); // Just fetch a few messages + CountDownLatch shutdownLatchStartpointTimestamp = new CountDownLatch(1); + + TestTaskApplication testTaskApplicationStartpointTimestamp = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic2, outputKafkaTopic, processedMessagesLatchStartpointTimestamp, + shutdownLatchStartpointTimestamp, Optional.of(processedCallback)); + + ApplicationRunner appRunner = + ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointTimestamp, applicationConfig2); executeRun(appRunner, applicationConfig2); assertTrue(processedMessagesLatchStartpointTimestamp.await(1, TimeUnit.MINUTES)); @@ -286,9 +285,51 @@ public void testStartpoints() throws InterruptedException { ime.getEventTime(), startpointTimestamp.getTimestampOffset()); assertTrue(assertMsg, ime.getEventTime() >= startpointTimestamp.getTimestampOffset()); } + } + + @Test + public void testStartpointOldest() throws InterruptedException { + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic3, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[2], Duration.ofMillis(2)); - // Startpoint oldest - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointOldest, applicationConfig3); + ConcurrentHashMap recvEventsInputStartpointOldest = new ConcurrentHashMap<>(); + + CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); + coordinatorStreamStore.init(); + StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); + startpointManager.start(); + + StartpointOldest startpointOldest = new StartpointOldest(); + writeStartpoints(startpointManager, inputKafkaTopic3, ZK_TEST_PARTITION_COUNT, startpointOldest); + + startpointManager.stop(); + coordinatorStreamStore.close(); + + TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + try { + String streamName = ime.getSystemStreamPartition().getStream(); + TestKafkaEvent testKafkaEvent = + TestKafkaEvent.fromString((String) ime.getMessage()); + String eventIndex = testKafkaEvent.getEventData(); + if (inputKafkaTopic3.equals(streamName)) { + recvEventsInputStartpointOldest.put(eventIndex, ime); + } else { + throw new RuntimeException("Unexpected input stream: " + streamName); + } + callback.complete(); + } catch (Exception ex) { + callback.failure(ex); + } + }; + + CountDownLatch processedMessagesLatchStartpointOldest = new CountDownLatch(NUM_KAFKA_EVENTS); // Fetch all since consuming from oldest + CountDownLatch shutdownLatchStartpointOldest = new CountDownLatch(1); + + TestTaskApplication testTaskApplicationStartpointOldest = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic3, outputKafkaTopic, processedMessagesLatchStartpointOldest, + shutdownLatchStartpointOldest, Optional.of(processedCallback)); + + ApplicationRunner appRunner = + ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointOldest, applicationConfig3); executeRun(appRunner, applicationConfig3); assertTrue(processedMessagesLatchStartpointOldest.await(1, TimeUnit.MINUTES)); @@ -299,9 +340,53 @@ public void testStartpoints() throws InterruptedException { assertTrue(shutdownLatchStartpointOldest.await(1, TimeUnit.MINUTES)); assertEquals("Expecting to have processed all the events", NUM_KAFKA_EVENTS, recvEventsInputStartpointOldest.size()); + } + + @Test + public void testStartpointUpcoming() throws InterruptedException { + publishKafkaEventsWithDelayPerEvent(inputKafkaTopic4, 0, NUM_KAFKA_EVENTS, PROCESSOR_IDS[3], Duration.ofMillis(2)); + + ConcurrentHashMap recvEventsInputStartpointUpcoming = new ConcurrentHashMap<>(); + + CoordinatorStreamStore coordinatorStreamStore = createCoordinatorStreamStore(applicationConfig1); + coordinatorStreamStore.init(); + StartpointManager startpointManager = new StartpointManager(coordinatorStreamStore); + startpointManager.start(); + + StartpointUpcoming startpointUpcoming = new StartpointUpcoming(); + writeStartpoints(startpointManager, inputKafkaTopic4, ZK_TEST_PARTITION_COUNT, startpointUpcoming); + + startpointManager.stop(); + coordinatorStreamStore.close(); + + TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + try { + String streamName = ime.getSystemStreamPartition().getStream(); + TestKafkaEvent testKafkaEvent = + TestKafkaEvent.fromString((String) ime.getMessage()); + String eventIndex = testKafkaEvent.getEventData(); + if (inputKafkaTopic4.equals(streamName)) { + recvEventsInputStartpointUpcoming.put(eventIndex, ime); + } else { + throw new RuntimeException("Unexpected input stream: " + streamName); + } + callback.complete(); + } catch (Exception ex) { + callback.failure(ex); + } + }; + + + CountDownLatch processedMessagesLatchStartpointUpcoming = new CountDownLatch(5); // Expecting none, so just attempt a small number of fetches. + CountDownLatch shutdownLatchStartpointUpcoming = new CountDownLatch(1); + + TestTaskApplication testTaskApplicationStartpointUpcoming = + new TestTaskApplication(TEST_SYSTEM, inputKafkaTopic4, outputKafkaTopic, processedMessagesLatchStartpointUpcoming, + shutdownLatchStartpointUpcoming, Optional.of(processedCallback)); // Startpoint upcoming - appRunner = ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointUpcoming, applicationConfig4); + ApplicationRunner appRunner = + ApplicationRunners.getApplicationRunner(testTaskApplicationStartpointUpcoming, applicationConfig4); executeRun(appRunner, applicationConfig4); assertFalse("Expecting to timeout and not process any old messages.", processedMessagesLatchStartpointUpcoming.await(15, TimeUnit.SECONDS)); @@ -314,8 +399,8 @@ public void testStartpoints() throws InterruptedException { } // Sends each event synchronously and returns map of event index to event metadata of the sent record/event. Adds the specified delay between sends. - private TreeMap publishKafkaEventsWithDelayPerEvent(String topic, int startIndex, int endIndex, String streamProcessorId, Duration delay) { - TreeMap eventsMetadata = new TreeMap<>(); + private Map publishKafkaEventsWithDelayPerEvent(String topic, int startIndex, int endIndex, String streamProcessorId, Duration delay) { + HashMap eventsMetadata = new HashMap<>(); for (int eventIndex = startIndex; eventIndex < endIndex; eventIndex++) { try { LOGGER.info("Publish kafka event with index : {} for stream processor: {}.", eventIndex, streamProcessorId); From 776932a38d3a7c6b3208ab9ee519671c0820e7fd Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Wed, 2 Oct 2019 12:37:08 -0700 Subject: [PATCH 5/6] Address comments from @mynameborat --- .../test/processor/TestTaskApplication.java | 6 +++--- ...pointTestHarness.java => TestStartpoint.java} | 16 +++++++++------- 2 files changed, 12 insertions(+), 10 deletions(-) rename samza-test/src/test/java/org/apache/samza/test/startpoint/{StartpointTestHarness.java => TestStartpoint.java} (95%) diff --git a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java index feba56e0a4..fabbbc95e2 100644 --- a/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java +++ b/samza-test/src/test/java/org/apache/samza/test/processor/TestTaskApplication.java @@ -47,7 +47,7 @@ public class TestTaskApplication implements TaskApplication { private final String outputTopic; private final CountDownLatch shutdownLatch; private final CountDownLatch processedMessageLatch; - private final Optional processCallback; + private final Optional processCallback; /** * A test TaskApplication to use in test harnesses. @@ -72,7 +72,7 @@ public TestTaskApplication(String systemName, String inputTopic, String outputTo * @param processCallback optional callback called per message processed. */ public TestTaskApplication(String systemName, String inputTopic, String outputTopic, - CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch, Optional processCallback) { + CountDownLatch processedMessageLatch, CountDownLatch shutdownLatch, Optional processCallback) { this.systemName = systemName; this.inputTopic = inputTopic; this.outputTopic = outputTopic; @@ -108,7 +108,7 @@ public void describe(TaskApplicationDescriptor appDescriptor) { .withTaskFactory((AsyncStreamTaskFactory) () -> new TestTaskImpl()); } - public interface TaskApplicationCallback { + public interface TaskApplicationProcessCallback { void onMessage(IncomingMessageEnvelope m, TaskCallback callback); } } diff --git a/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java b/samza-test/src/test/java/org/apache/samza/test/startpoint/TestStartpoint.java similarity index 95% rename from samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java rename to samza-test/src/test/java/org/apache/samza/test/startpoint/TestStartpoint.java index c064b7bdfd..d71028d07f 100644 --- a/samza-test/src/test/java/org/apache/samza/test/startpoint/StartpointTestHarness.java +++ b/samza-test/src/test/java/org/apache/samza/test/startpoint/TestStartpoint.java @@ -80,9 +80,9 @@ import static org.junit.Assert.assertTrue; -public class StartpointTestHarness extends IntegrationTestHarness { +public class TestStartpoint extends IntegrationTestHarness { - private static final Logger LOGGER = LoggerFactory.getLogger(StartpointTestHarness.class); + private static final Logger LOGGER = LoggerFactory.getLogger(TestStartpoint.class); private static final int NUM_KAFKA_EVENTS = 300; private static final int ZK_TEST_PARTITION_COUNT = 5; @@ -108,8 +108,10 @@ public class StartpointTestHarness extends IntegrationTestHarness { private String testStreamAppName; private String testStreamAppId; + // Tests typically take 20-30 seconds each due to timers such as debounce time, embedded ZK and Kafka initialization, etc... + // Capping the timeout here with some buffer at 60 seconds to account for possible noisy neighbors while running the tests. @Rule - public Timeout testTimeOutInMillis = new Timeout(150000, TimeUnit.MILLISECONDS); + public Timeout testTimeOutInMillis = new Timeout(60000, TimeUnit.MILLISECONDS); @Rule public final ExpectedException expectedException = ExpectedException.none(); @@ -182,7 +184,7 @@ public void testStartpointSpecific() throws InterruptedException { startpointManager.stop(); coordinatorStreamStore.close(); - TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + TestTaskApplication.TaskApplicationProcessCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { try { String streamName = ime.getSystemStreamPartition().getStream(); TestKafkaEvent testKafkaEvent = TestKafkaEvent.fromString((String) ime.getMessage()); @@ -243,7 +245,7 @@ public void testStartpointTimestamp() throws InterruptedException { startpointManager.stop(); coordinatorStreamStore.close(); - TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + TestTaskApplication.TaskApplicationProcessCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { try { String streamName = ime.getSystemStreamPartition().getStream(); TestKafkaEvent testKafkaEvent = @@ -304,7 +306,7 @@ public void testStartpointOldest() throws InterruptedException { startpointManager.stop(); coordinatorStreamStore.close(); - TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + TestTaskApplication.TaskApplicationProcessCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { try { String streamName = ime.getSystemStreamPartition().getStream(); TestKafkaEvent testKafkaEvent = @@ -359,7 +361,7 @@ public void testStartpointUpcoming() throws InterruptedException { startpointManager.stop(); coordinatorStreamStore.close(); - TestTaskApplication.TaskApplicationCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { + TestTaskApplication.TaskApplicationProcessCallback processedCallback = (IncomingMessageEnvelope ime, TaskCallback callback) -> { try { String streamName = ime.getSystemStreamPartition().getStream(); TestKafkaEvent testKafkaEvent = From c4291ad2312a93fe4f3d24e433e41e06543d963c Mon Sep 17 00:00:00 2001 From: Daniel Nishimura Date: Wed, 2 Oct 2019 13:14:16 -0700 Subject: [PATCH 6/6] Trigger build