implements Runnable {
- private static final Log LOG = LogFactory.getLog(KinesisConnectorExecutorBase.class);
-
+ private static final Logger LOG = Logger.getLogger(KinesisConnectorExecutorBase.class.getName());
+
// Amazon Kinesis Client Library worker to process records
protected Worker worker;
diff --git a/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessor.java b/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessor.java
new file mode 100644
index 0000000..5c291ce
--- /dev/null
+++ b/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessor.java
@@ -0,0 +1,213 @@
+/*
+ * Copyright 2013-2014 Amazon.com, Inc. or its affiliates. All Rights Reserved.
+ *
+ * Licensed under the Amazon Software License (the "License").
+ * You may not use this file except in compliance with the License.
+ * A copy of the License is located at
+ *
+ * http://aws.amazon.com/asl/
+ *
+ * or in the "license" file accompanying this file. This file 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 com.sumologic.kinesis;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.List;
+
+import org.apache.log4j.Logger;
+
+import com.amazonaws.services.kinesis.clientlibrary.exceptions.InvalidStateException;
+import com.amazonaws.services.kinesis.clientlibrary.exceptions.KinesisClientLibDependencyException;
+import com.amazonaws.services.kinesis.clientlibrary.exceptions.ShutdownException;
+import com.amazonaws.services.kinesis.clientlibrary.exceptions.ThrottlingException;
+import com.amazonaws.services.kinesis.clientlibrary.interfaces.IRecordProcessor;
+import com.amazonaws.services.kinesis.clientlibrary.interfaces.IRecordProcessorCheckpointer;
+import com.amazonaws.services.kinesis.clientlibrary.types.ShutdownReason;
+import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
+import com.amazonaws.services.kinesis.connectors.UnmodifiableBuffer;
+import com.amazonaws.services.kinesis.connectors.interfaces.IBuffer;
+import com.amazonaws.services.kinesis.connectors.interfaces.ICollectionTransformer;
+import com.amazonaws.services.kinesis.connectors.interfaces.IEmitter;
+import com.amazonaws.services.kinesis.connectors.interfaces.IFilter;
+import com.amazonaws.services.kinesis.connectors.interfaces.ITransformer;
+import com.amazonaws.services.kinesis.connectors.interfaces.ITransformerBase;
+import com.amazonaws.services.kinesis.model.Record;
+
+/**
+ * This is the base class for any KinesisConnector. It is configured by a constructor that takes in
+ * as parameters implementations of the IBuffer, ITransformer, and IEmitter dependencies defined in
+ * a IKinesisConnectorPipeline. It is typed to match the class that records are transformed into for
+ * filtering and manipulation. This class is produced by a KinesisConnectorRecordProcessorFactory.
+ *
+ * When a Worker calls processRecords() on this class, the pipeline is used in the following way:
+ *
+ * - Records are transformed into the corresponding data model (parameter type T) via the ITransformer.
+ * - Transformed records are passed to the IBuffer.consumeRecord() method, which may optionally filter based on the
+ * IFilter in the pipeline.
+ * - When the buffer is full (IBuffer.shouldFlush() returns true), records are transformed with the ITransformer to
+ * the output type (parameter type U) and a call is made to IEmitter.emit(). IEmitter.emit() returning an empty list is
+ * considered a success, so the record processor will checkpoint and emit will not be retried. Non-empty return values
+ * will result in additional calls to emit with failed records as the unprocessed list until the retry limit is reached.
+ * Upon exceeding the retry limit or an exception being thrown, the IEmitter.fail() method will be called with the
+ * unprocessed records.
+ * - When the shutdown() method of this class is invoked, a call is made to the IEmitter.shutdown() method which
+ * should close any existing client connections.
+ *
+ *
+ */
+public class KinesisConnectorRecordProcessor implements IRecordProcessor {
+
+ private final IEmitter emitter;
+ private final ITransformerBase transformer;
+ private final IFilter filter;
+ private final IBuffer buffer;
+ private final int retryLimit;
+ private final long backoffInterval;
+ private boolean isShutdown = false;
+
+ private static final Logger LOG = Logger.getLogger(KinesisConnectorRecordProcessor.class.getName());
+
+ private String shardId;
+
+ public KinesisConnectorRecordProcessor(IBuffer buffer,
+ IFilter filter,
+ IEmitter emitter,
+ ITransformerBase transformer,
+ KinesisConnectorConfiguration configuration) {
+ if (buffer == null || filter == null || emitter == null || transformer == null) {
+ throw new IllegalArgumentException("buffer, filter, emitter, and transformer must not be null");
+ }
+ this.buffer = buffer;
+ this.filter = filter;
+ this.emitter = emitter;
+ this.transformer = transformer;
+ // Limit must be greater than zero
+ if (configuration.RETRY_LIMIT <= 0) {
+ retryLimit = 1;
+ } else {
+ retryLimit = configuration.RETRY_LIMIT;
+ }
+ this.backoffInterval = configuration.BACKOFF_INTERVAL;
+ }
+
+ @Override
+ public void initialize(String shardId) {
+ this.shardId = shardId;
+ }
+
+ @Override
+ public void processRecords(List records, IRecordProcessorCheckpointer checkpointer) {
+ // Note: This method will be called even for empty record lists. This is needed for checking the buffer time
+ // threshold.
+ if (isShutdown) {
+ LOG.warn("processRecords called on shutdown record processor for shardId: " + shardId);
+ return;
+ }
+ if (shardId == null) {
+ throw new IllegalStateException("Record processor not initialized");
+ }
+
+ // Transform each Amazon Kinesis Record and add the result to the buffer
+ for (Record record : records) {
+ try {
+ if (transformer instanceof ITransformer) {
+ ITransformer singleTransformer = (ITransformer) transformer;
+ filterAndBufferRecord(singleTransformer.toClass(record), record);
+ } else if (transformer instanceof ICollectionTransformer) {
+ ICollectionTransformer listTransformer = (ICollectionTransformer) transformer;
+ Collection transformedRecords = listTransformer.toClass(record);
+ for (T transformedRecord : transformedRecords) {
+ filterAndBufferRecord(transformedRecord, record);
+ }
+ } else {
+ throw new RuntimeException("Transformer must implement ITransformer or ICollectionTransformer");
+ }
+ } catch (IOException e) {
+ LOG.error(e);
+ }
+ }
+
+ if (buffer.shouldFlush()) {
+ List emitItems = transformToOutput(buffer.getRecords());
+ emit(checkpointer, emitItems);
+ }
+ }
+
+ private void filterAndBufferRecord(T transformedRecord, Record record) {
+ if (filter.keepRecord(transformedRecord)) {
+ buffer.consumeRecord(transformedRecord, record.getData().array().length, record.getSequenceNumber());
+ }
+ }
+
+ private List transformToOutput(List items) {
+ List emitItems = new ArrayList();
+ for (T item : items) {
+ try {
+ emitItems.add(transformer.fromClass(item));
+ } catch (IOException e) {
+ LOG.error("Failed to transform record " + item + " to output type", e);
+ }
+ }
+ return emitItems;
+ }
+
+ private void emit(IRecordProcessorCheckpointer checkpointer, List emitItems) {
+ List unprocessed = new ArrayList(emitItems);
+ try {
+ for (int numTries = 0; numTries < retryLimit; numTries++) {
+ unprocessed = emitter.emit(new UnmodifiableBuffer(buffer, unprocessed));
+ if (unprocessed.isEmpty()) {
+ break;
+ }
+ try {
+ Thread.sleep(backoffInterval);
+ } catch (InterruptedException e) {
+ }
+ }
+ if (!unprocessed.isEmpty()) {
+ emitter.fail(unprocessed);
+ }
+ final String lastSequenceNumberProcessed = buffer.getLastSequenceNumber();
+ buffer.clear();
+ // checkpoint once all the records have been consumed
+ if (lastSequenceNumberProcessed != null && unprocessed.isEmpty()) {
+ checkpointer.checkpoint(lastSequenceNumberProcessed);
+ }
+ } catch (IOException | KinesisClientLibDependencyException | InvalidStateException | ThrottlingException
+ | ShutdownException e) {
+ LOG.error(e);
+ emitter.fail(unprocessed);
+ }
+ }
+
+ @Override
+ public void shutdown(IRecordProcessorCheckpointer checkpointer, ShutdownReason reason) {
+ LOG.info("Shutting down record processor with shardId: " + shardId + " with reason " + reason);
+ if (isShutdown) {
+ LOG.warn("Record processor for shardId: " + shardId + " has been shutdown multiple times.");
+ return;
+ }
+ switch (reason) {
+ case TERMINATE:
+ emit(checkpointer, transformToOutput(buffer.getRecords()));
+ try {
+ checkpointer.checkpoint();
+ } catch (KinesisClientLibDependencyException | InvalidStateException | ThrottlingException | ShutdownException e) {
+ LOG.error(e);
+ }
+ break;
+ case ZOMBIE:
+ break;
+ default:
+ throw new IllegalStateException("invalid shutdown reason");
+ }
+ emitter.shutdown();
+ isShutdown = true;
+ }
+
+}
\ No newline at end of file
diff --git a/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessorFactory.java b/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessorFactory.java
new file mode 100644
index 0000000..2144e62
--- /dev/null
+++ b/src/main/java/com/sumologic/kinesis/KinesisConnectorRecordProcessorFactory.java
@@ -0,0 +1,43 @@
+package com.sumologic.kinesis;
+
+import com.amazonaws.services.kinesis.clientlibrary.interfaces.IRecordProcessor;
+import com.amazonaws.services.kinesis.clientlibrary.interfaces.IRecordProcessorFactory;
+import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
+import com.amazonaws.services.kinesis.connectors.interfaces.IBuffer;
+import com.amazonaws.services.kinesis.connectors.interfaces.IEmitter;
+import com.amazonaws.services.kinesis.connectors.interfaces.IFilter;
+import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
+import com.amazonaws.services.kinesis.connectors.interfaces.ITransformerBase;
+
+/**
+ * This class is used to generate KinesisConnectorRecordProcessors that operate using the user's
+ * implemented classes. The createProcessor() method sets the dependencies of the
+ * KinesisConnectorRecordProcessor that are specified in the KinesisConnectorPipeline argument,
+ * which accesses instances of the users implementations.
+ */
+public class KinesisConnectorRecordProcessorFactory implements IRecordProcessorFactory {
+
+ private IKinesisConnectorPipeline pipeline;
+ private KinesisConnectorConfiguration configuration;
+
+ public KinesisConnectorRecordProcessorFactory(IKinesisConnectorPipeline pipeline,
+ KinesisConnectorConfiguration configuration) {
+ this.configuration = configuration;
+ this.pipeline = pipeline;
+ }
+
+ @Override
+ public IRecordProcessor createProcessor() {
+ try {
+ IBuffer buffer = pipeline.getBuffer(configuration);
+ IEmitter emitter = pipeline.getEmitter(configuration);
+ ITransformerBase transformer = pipeline.getTransformer(configuration);
+ IFilter filter = pipeline.getFilter(configuration);
+ KinesisConnectorRecordProcessor processor =
+ new KinesisConnectorRecordProcessor(buffer, filter, emitter, transformer, configuration);
+ return processor;
+ } catch (Throwable t) {
+ throw new RuntimeException(t);
+ }
+ }
+}
\ No newline at end of file
diff --git a/src/main/java/com/sumologic/kinesis/StreamSource.java b/src/main/java/com/sumologic/kinesis/StreamSource.java
index 0290b5d..ab74b37 100644
--- a/src/main/java/com/sumologic/kinesis/StreamSource.java
+++ b/src/main/java/com/sumologic/kinesis/StreamSource.java
@@ -6,12 +6,10 @@
import java.io.InputStreamReader;
import java.nio.ByteBuffer;
-import org.apache.commons.logging.Log;
-import org.apache.commons.logging.LogFactory;
+import org.apache.log4j.Logger;
-import com.sumologic.client.SimpleKinesisMessageModel;
+import com.sumologic.client.model.SimpleKinesisMessageModel;
import com.sumologic.kinesis.utils.KinesisUtils;
-
import com.amazonaws.auth.AWSCredentialsProvider;
import com.amazonaws.regions.RegionUtils;
import com.amazonaws.services.kinesis.AmazonKinesisClient;
@@ -25,7 +23,7 @@
* stream defined in the KinesisConnectorConfiguration.
*/
public class StreamSource implements Runnable {
- private static Log LOG = LogFactory.getLog(StreamSource.class);
+ private static final Logger LOG = Logger.getLogger(StreamSource.class.getName());
protected AmazonKinesisClient kinesisClient;
protected KinesisConnectorConfiguration config;
protected final String inputFile;
@@ -106,7 +104,8 @@ protected void processInputStream(InputStream inputStream, int iteration) throws
String line;
int lines = 0;
while ((line = br.readLine()) != null) {
- SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
+ SimpleKinesisMessageModel kinesisMessageModel = new SimpleKinesisMessageModel(line);
+ //SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
PutRecordRequest putRecordRequest = new PutRecordRequest();
putRecordRequest.setStreamName(config.KINESIS_INPUT_STREAM);
diff --git a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
index 1588a77..dab6ad0 100644
--- a/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
+++ b/src/test/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformerTest.java
@@ -7,8 +7,8 @@
import org.junit.Ignore;
import com.amazonaws.services.kinesis.model.Record;
-import com.sumologic.client.CloudWatchMessageModelSumologicTransformer;
-import com.sumologic.client.SimpleKinesisMessageModel;
+import com.sumologic.client.model.CloudWatchLogsMessageModel;
+import com.sumologic.client.model.SimpleKinesisMessageModel;
import java.io.IOException;
import java.nio.charset.Charset;
@@ -61,7 +61,7 @@ public void theTransformerShouldSucceedWhenTransformingAProperJSON() {
+ "}"
+"";
- byte[] compressData = SumologicSender.compressGzip(jsonData);
+ byte[] compressData = SumologicKinesisUtils.compressGzip(jsonData);
ByteBuffer bufferedData = null;
try {
@@ -97,7 +97,7 @@ public void theTransformerShouldFailWhenTransformingAJSONWithTrailingCommas() {
+ "}"
+"";
- byte[] compressData = SumologicSender.compressGzip(jsonData);
+ byte[] compressData = SumologicKinesisUtils.compressGzip(jsonData);
ByteBuffer bufferedData = null;
try {
@@ -139,7 +139,7 @@ public void theTransfomerShouldSeparateBatchesOfLogs() {
+ "}"
+"";
- byte[] compressData = SumologicSender.compressGzip(jsonData);
+ byte[] compressData = SumologicKinesisUtils.compressGzip(jsonData);
ByteBuffer bufferedData = null;
try {
diff --git a/src/test/java/com/sumologic/client/SumologicKinesisUtilsTest.java b/src/test/java/com/sumologic/client/SumologicKinesisUtilsTest.java
new file mode 100644
index 0000000..5982273
--- /dev/null
+++ b/src/test/java/com/sumologic/client/SumologicKinesisUtilsTest.java
@@ -0,0 +1,56 @@
+package com.sumologic.client;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+public class SumologicKinesisUtilsTest {
+ @Test
+ public void compressDecompressGzipTest() {
+ String data = "a string of characters";
+
+ byte[] compressData = SumologicKinesisUtils.compressGzip(data);
+ String result = SumologicKinesisUtils.decompressGzip(compressData);
+
+ Assert.assertTrue(data.equals(result));
+ }
+
+ @Test
+ public void properJSONVerificationShouldReturnTrue() {
+ String jsonData = ""
+ +"{"
+ + "\"logEvents\": [{"
+ + "\"id\": \"3889492387492837492374982374897239847289374892\","
+ + "\"message\": \"1 23423532532 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 ACCEPT OK\","
+ + "\"timestamp\": \"2342342342300\""
+ + "}],"
+ + "\"logGroup\": \"MyFirstVPC\","
+ + "\"logStream\": \"eni-ac6a7de4-all\","
+ + "\"messageType\": \"DATA_MESSAGE\","
+ + "\"owner\": \"2342352352\","
+ + "\"subscriptionFilters\": [\"MyFirstVPC\"]"
+ + "}"
+ +"";
+
+ Assert.assertTrue(SumologicKinesisUtils.verifyJSON(jsonData));
+ }
+
+ @Test
+ public void malformedJSONVerificationShouldReturnTrue() {
+ String jsonData = ""
+ +"{"
+ + "\"logEvents\": [{"
+ + "\"id\": \"3889492387492837492374982374897239847289374892\","
+ + "\"message\": \"1 23423532532 eni-ac9342k3492 10.1.1.75 66.175.209.17 123 123 17 1 76 1437755534 1437755549 ACCEPT OK\","
+ + "\"timestamp\": \"2342342342300\""
+ + "}],"
+ + "\"logGroup\": \"MyFirstVPC\","
+ + "\"logStream\": \"eni-ac6a7de4-all\","
+ + "\"messageType\": \"DATA_MESSAGE\","
+ + "\"owner\": \"2342352352\","
+ + "\"subscriptionFilters\": [\"MyFirstVPC\"],"
+ + "}"
+ +"";
+
+ Assert.assertFalse(SumologicKinesisUtils.verifyJSON(jsonData));
+ }
+}
\ No newline at end of file
diff --git a/src/test/java/com/sumologic/client/SumologicSenderTest.java b/src/test/java/com/sumologic/client/SumologicSenderTest.java
index 1d9c54d..215ac3f 100644
--- a/src/test/java/com/sumologic/client/SumologicSenderTest.java
+++ b/src/test/java/com/sumologic/client/SumologicSenderTest.java
@@ -16,7 +16,6 @@
import com.github.tomakehurst.wiremock.client.WireMock;
import com.github.tomakehurst.wiremock.junit.WireMockRule;
-import com.sumologic.client.CloudWatchMessageModelSumologicTransformer;
import com.sumologic.client.SumologicSender;
import com.sumologic.client.implementations.SumologicEmitter;
@@ -68,18 +67,7 @@ public void theSenderShouldReturnTrueOnSuccess () {
}
}
- @Test
- public void decompressGzipTest() {
- String url = MOCKED_HOST + MOCKED_COLLECTION;
-
- String data = "a string of characters";
-
- byte[] compressData = SumologicSender.compressGzip(data);
- String result = CloudWatchMessageModelSumologicTransformer.decompressGzip(compressData);
-
- Assert.assertTrue(data.equals(result));
- }
-
+
private void mockEmitMessages () {
WireMock.stubFor(WireMock.post(WireMock.urlMatching(MOCKED_COLLECTION))
.willReturn(WireMock.aResponse()