Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,10 @@

import java.io.IOException;
import java.io.ObjectInputStream;
import java.io.Serializable;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import org.apache.samza.application.descriptors.StreamApplicationDescriptor;
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;
Expand All @@ -37,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;


/**
Expand Down Expand Up @@ -114,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 = message.split("|");
return new TestKafkaEvent(messageComponents[0], messageComponents[1]);
}
}

public static StreamApplication getInstance(
String systemName,
List<String> inputTopics,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -46,14 +47,38 @@ public class TestTaskApplication implements TaskApplication {
private final String outputTopic;
private final CountDownLatch shutdownLatch;
private final CountDownLatch processedMessageLatch;
private final Optional<TaskApplicationProcessCallback> 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) {
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<TaskApplicationProcessCallback> 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 {
Expand All @@ -62,7 +87,8 @@ 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.
processCallback.ifPresent(pcb -> pcb.onMessage(envelope, callback));
}

@Override
Expand All @@ -81,4 +107,8 @@ public void describe(TaskApplicationDescriptor appDescriptor) {
.withOutputStream(outputDescriptor)
.withTaskFactory((AsyncStreamTaskFactory) () -> new TestTaskImpl());
}

public interface TaskApplicationProcessCallback {
void onMessage(IncomingMessageEnvelope m, TaskCallback callback);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,29 +27,29 @@
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.UUID;
import java.util.concurrent.CountDownLatch;
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.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;
Expand All @@ -68,12 +68,13 @@
import org.apache.samza.system.SystemStreamPartition;
import org.apache.samza.test.StandaloneTestUtils;
import org.apache.samza.test.harness.IntegrationTestHarness;
import org.apache.samza.test.util.TestKafkaEvent;
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;
Expand All @@ -83,7 +84,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;


/**
Expand Down Expand Up @@ -184,7 +188,7 @@ 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()));
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);
Expand Down Expand Up @@ -540,7 +544,7 @@ public void testRollingUpgradeOfStreamApplicationsShouldGenerateSameJobModel() t
configMap.put(JobConfig.PROCESSOR_ID, PROCESSOR_IDS[1]);
Config applicationConfig2 = new MapConfig(configMap);

List<TestStreamApplication.TestKafkaEvent> messagesProcessed = new ArrayList<>();
List<TestKafkaEvent> messagesProcessed = new ArrayList<>();
TestStreamApplication.StreamApplicationCallback streamApplicationCallback = messagesProcessed::add;

// Create StreamApplication from configuration.
Expand Down Expand Up @@ -570,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();
Expand Down
Loading