Skip to content

Commit 2403a31

Browse files
author
Juan Pablo Diaz-Vaz
committed
Merged in SUMOK-18 (pull request #6)
Sumok 18
2 parents 5e76d15 + 4c8b07d commit 2403a31

12 files changed

Lines changed: 145 additions & 94 deletions

src/main/java/com/mcplusa/kinesis/BatchedStreamSource.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@
2828
import org.apache.commons.logging.Log;
2929
import org.apache.commons.logging.LogFactory;
3030

31-
import com.mcplusa.sumologic.KinesisMessageModel;
31+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
3232

3333
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
3434
import com.amazonaws.services.kinesis.model.PutRecordRequest;
@@ -41,15 +41,15 @@ public class BatchedStreamSource extends StreamSource {
4141
private static Log LOG = LogFactory.getLog(BatchedStreamSource.class);
4242

4343
private static int NUM_BYTES_PER_PUT_REQUEST = 50000;
44-
List<KinesisMessageModel> buffer;
44+
List<SimpleKinesisMessageModel> buffer;
4545

4646
public BatchedStreamSource(KinesisConnectorConfiguration config, String inputFile) {
4747
this(config, inputFile, false);
4848
}
4949

5050
public BatchedStreamSource(KinesisConnectorConfiguration config, String inputFile, boolean loopOverStreamSource) {
5151
super(config, inputFile, loopOverStreamSource);
52-
buffer = new ArrayList<KinesisMessageModel>();
52+
buffer = new ArrayList<SimpleKinesisMessageModel>();
5353
}
5454

5555
@Override
@@ -59,14 +59,14 @@ protected void processInputStream(InputStream inputStream, int iteration) throws
5959
int lines = 0;
6060

6161
while ((line = br.readLine()) != null) {
62-
KinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, KinesisMessageModel.class);
62+
SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
6363
buffer.add(kinesisMessageModel);
6464
if (numBytesInBuffer() > NUM_BYTES_PER_PUT_REQUEST) {
6565
/*
6666
* We need to remove the last record to ensure this data blob is accepted by the Amazon Kinesis
6767
* client which restricts the data blob to be less than 50 KB.
6868
*/
69-
KinesisMessageModel lastRecord = buffer.remove(buffer.size() - 1);
69+
SimpleKinesisMessageModel lastRecord = buffer.remove(buffer.size() - 1);
7070
flushBuffer();
7171
/*
7272
* We add it back so it will be part of the next batch.

src/main/java/com/mcplusa/kinesis/KinesisConnectorExecutor.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,9 +38,11 @@ public abstract class KinesisConnectorExecutor<T, U> extends KinesisConnectorExe
3838
// Create Stream Source constants
3939
private static final String CREATE_STREAM_SOURCE = "createStreamSource";
4040
private static final String LOOP_OVER_STREAM_SOURCE = "loopOverStreamSource";
41+
private static final String INPUT_STREAM_FILE = "inputStreamFile";
42+
4143
private static final boolean DEFAULT_CREATE_STREAM_SOURCE = false;
4244
private static final boolean DEFAULT_LOOP_OVER_STREAM_SOURCE = false;
43-
private static final String INPUT_STREAM_FILE = "inputStreamFile";
45+
4446

4547
// Class variables
4648
protected final KinesisConnectorForSumologicConfiguration config;

src/main/java/com/mcplusa/kinesis/StreamSource.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
import org.apache.commons.logging.Log;
2424
import org.apache.commons.logging.LogFactory;
2525

26-
import com.mcplusa.sumologic.KinesisMessageModel;
26+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
2727
import com.mcplusa.kinesis.utils.KinesisUtils;
2828

2929
import com.amazonaws.auth.AWSCredentialsProvider;
@@ -120,7 +120,7 @@ protected void processInputStream(InputStream inputStream, int iteration) throws
120120
String line;
121121
int lines = 0;
122122
while ((line = br.readLine()) != null) {
123-
KinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, KinesisMessageModel.class);
123+
SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
124124

125125
PutRecordRequest putRecordRequest = new PutRecordRequest();
126126
putRecordRequest.setStreamName(config.KINESIS_INPUT_STREAM);

src/main/java/com/mcplusa/sumologic/KinesisMessageModelSumologicTransformer.java renamed to src/main/java/com/mcplusa/sumologic/CloudWatchMessageModelSumologicTransformer.java

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44

55
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
66
import com.amazonaws.services.kinesis.model.Record;
7-
import com.mcplusa.sumologic.KinesisMessageModel;
7+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
88
import com.mcplusa.sumologic.implementations.SumologicEmitter;
99
import com.mcplusa.sumologic.implementations.SumologicTransformer;
1010

@@ -22,32 +22,32 @@
2222

2323

2424
/**
25-
* A custom transfomer for {@link KinesisMessageModel} records in JSON format. The output is in a format
25+
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
2626
* usable for insertions to Sumologic.
2727
*/
28-
public class KinesisMessageModelSumologicTransformer implements
29-
SumologicTransformer<KinesisMessageModel> {
28+
public class CloudWatchMessageModelSumologicTransformer implements
29+
SumologicTransformer<SimpleKinesisMessageModel> {
3030

31-
private static final Log LOG = LogFactory.getLog(KinesisMessageModelSumologicTransformer.class);
31+
private static final Log LOG = LogFactory.getLog(CloudWatchMessageModelSumologicTransformer.class);
3232

3333
/**
3434
* Creates a new KinesisMessageModelSumologicTransformer.
3535
*/
36-
public KinesisMessageModelSumologicTransformer() {
36+
public CloudWatchMessageModelSumologicTransformer() {
3737
super();
3838
}
3939

4040
@Override
41-
public String fromClass(KinesisMessageModel message) {
41+
public String fromClass(SimpleKinesisMessageModel message) {
4242
return message.toString();
4343
}
4444

4545
@Override
46-
public KinesisMessageModel toClass(Record record) throws IOException {
46+
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
4747
byte[] decodedRecord = record.getData().array();
4848
String stringifiedRecord = decompressGzip(decodedRecord);
4949

50-
return new KinesisMessageModel(stringifiedRecord);
50+
return new SimpleKinesisMessageModel(stringifiedRecord);
5151
}
5252

5353
public static String decompressGzip(byte[] compressedData) {
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package com.mcplusa.sumologic;
2+
3+
import java.io.IOException;
4+
5+
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
6+
import com.amazonaws.services.kinesis.model.Record;
7+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
8+
import com.mcplusa.sumologic.implementations.SumologicEmitter;
9+
import com.mcplusa.sumologic.implementations.SumologicTransformer;
10+
11+
import org.apache.commons.logging.Log;
12+
import org.apache.commons.logging.LogFactory;
13+
14+
import java.io.ByteArrayInputStream;
15+
import java.io.BufferedReader;
16+
import java.io.InputStreamReader;
17+
import java.util.zip.GZIPInputStream;
18+
import java.nio.charset.StandardCharsets;
19+
import java.util.Arrays;
20+
21+
import org.apache.commons.codec.binary.Base64;
22+
23+
24+
/**
25+
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
26+
* usable for insertions to Sumologic.
27+
*/
28+
public class DefaultKinesisMessageModelSumologicTransformer implements
29+
SumologicTransformer<SimpleKinesisMessageModel> {
30+
31+
private static final Log LOG = LogFactory.getLog(DefaultKinesisMessageModelSumologicTransformer.class);
32+
33+
/**
34+
* Creates a new KinesisMessageModelSumologicTransformer.
35+
*/
36+
public DefaultKinesisMessageModelSumologicTransformer() {
37+
super();
38+
}
39+
40+
@Override
41+
public String fromClass(SimpleKinesisMessageModel message) {
42+
return message.toString();
43+
}
44+
45+
@Override
46+
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
47+
byte[] decodedRecord = record.getData().array();
48+
String stringifiedRecord = new String(decodedRecord);
49+
50+
return new SimpleKinesisMessageModel(stringifiedRecord);
51+
}
52+
}

src/main/java/com/mcplusa/sumologic/KinesisConnectorForSumologicConfiguration.java

Lines changed: 9 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
import com.amazonaws.auth.AWSCredentialsProvider;
66
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
77

8-
98
/**
109
* This class contains constants used to configure AWS Services in Amazon Kinesis Connectors. The user
1110
* should use System properties to set their proper configuration. An instance of
@@ -14,25 +13,22 @@
1413
public class KinesisConnectorForSumologicConfiguration extends KinesisConnectorConfiguration {
1514
// Properties added for Sumologic
1615
public static final String PROP_SUMOLOGIC_URL = "sumologicUrl";
17-
public static final String PROP_SUMOLOGIC_USE_LOG4J = "useLog4j";
18-
public final String SUMOLOGIC_URL;
19-
public final boolean SUMOLOGIC_USE_LOG4J;
16+
public static final String PROP_TRANSFORMER_CLASS = "transformerClass";
2017

21-
public final boolean DEFAULT_SUMOLOGIC_USE_LOG4J = false;
18+
private static final String DEFAULT_SUMOLOGIC_URL = null;
19+
private static final String DEFAULT_TRANSFORMER_CLASS = null;
20+
21+
public final String SUMOLOGIC_URL;
22+
public final String TRANSFORMER_CLASS;
2223

2324
/**
2425
* Configure the connector application with any set of properties that are unique to the application. Any
2526
* unspecified property will be set to a default value.
2627
*/
2728
public KinesisConnectorForSumologicConfiguration(Properties properties, AWSCredentialsProvider credentialsProvider) {
2829
super(properties, credentialsProvider);
29-
SUMOLOGIC_URL = properties.getProperty(PROP_SUMOLOGIC_URL, null);
30-
SUMOLOGIC_USE_LOG4J = getBooleanProperty(PROP_SUMOLOGIC_USE_LOG4J,
31-
DEFAULT_SUMOLOGIC_USE_LOG4J, properties);
30+
31+
SUMOLOGIC_URL = properties.getProperty(PROP_SUMOLOGIC_URL, DEFAULT_SUMOLOGIC_URL);
32+
TRANSFORMER_CLASS = properties.getProperty(PROP_TRANSFORMER_CLASS, DEFAULT_TRANSFORMER_CLASS);
3233
}
33-
34-
private boolean getBooleanProperty(String property, boolean defaultValue, Properties properties) {
35-
String propertyValue = properties.getProperty(property, Boolean.toString(defaultValue));
36-
return Boolean.parseBoolean(propertyValue);
37-
}
3834
}

src/main/java/com/mcplusa/sumologic/KinesisMessageModel.java renamed to src/main/java/com/mcplusa/sumologic/SimpleKinesisMessageModel.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,11 @@
22

33
import java.io.Serializable;
44

5-
public class KinesisMessageModel implements Serializable {
5+
public class SimpleKinesisMessageModel implements Serializable {
66
private String data;
77
private int id;
88

9-
public KinesisMessageModel(String data) {
9+
public SimpleKinesisMessageModel(String data) {
1010
this.data = data;
1111
this.id = 1;
1212
}

src/main/java/com/mcplusa/sumologic/SumologicExecutor.java

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,10 @@
33
import com.amazonaws.services.kinesis.connectors.KinesisConnectorRecordProcessorFactory;
44

55
import com.mcplusa.kinesis.KinesisConnectorExecutor;
6-
import com.mcplusa.sumologic.KinesisMessageModel;
6+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
77
import com.mcplusa.sumologic.SumologicMessageModelPipeline;
88

9-
public class SumologicExecutor extends KinesisConnectorExecutor<KinesisMessageModel, String> {
9+
public class SumologicExecutor extends KinesisConnectorExecutor<SimpleKinesisMessageModel, String> {
1010

1111
private static String configFile = "SumologicConnector.properties";
1212

@@ -19,9 +19,9 @@ public SumologicExecutor(String configFile) {
1919
}
2020

2121
@Override
22-
public KinesisConnectorRecordProcessorFactory<KinesisMessageModel, String>
22+
public KinesisConnectorRecordProcessorFactory<SimpleKinesisMessageModel, String>
2323
getKinesisConnectorRecordProcessorFactory() {
24-
return new KinesisConnectorRecordProcessorFactory<KinesisMessageModel, String>
24+
return new KinesisConnectorRecordProcessorFactory<SimpleKinesisMessageModel, String>
2525
(new SumologicMessageModelPipeline(),config);
2626
}
2727

@@ -30,7 +30,7 @@ public SumologicExecutor(String configFile) {
3030
* @param args
3131
*/
3232
public static void main(String[] args) {
33-
KinesisConnectorExecutor<KinesisMessageModel, String> sumologicExecutor =
33+
KinesisConnectorExecutor<SimpleKinesisMessageModel, String> sumologicExecutor =
3434
new SumologicExecutor(configFile);
3535
sumologicExecutor.run();
3636
}
Lines changed: 34 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
package com.mcplusa.sumologic;
22

3-
import com.mcplusa.sumologic.KinesisMessageModel;
4-
import com.mcplusa.sumologic.KinesisMessageModelSumologicTransformer;
5-
import com.mcplusa.sumologic.implementations.SumologicEmitter;
3+
import org.apache.commons.logging.Log;
4+
import org.apache.commons.logging.LogFactory;
65

6+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
7+
import com.mcplusa.sumologic.CloudWatchMessageModelSumologicTransformer;
8+
import com.mcplusa.sumologic.implementations.SumologicEmitter;
79
import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
810
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
911
import com.amazonaws.services.kinesis.connectors.impl.BasicMemoryBuffer;
@@ -20,32 +22,53 @@
2022
* <ul>
2123
* <li>{@link SumologicEmitter}</li>
2224
* <li>{@link BasicMemoryBuffer}</li>
23-
* <li>{@link KinesisMessageModelSumologicTransformer}</li>
25+
* <li>{@link CloudWatchMessageModelSumologicTransformer}</li>
2426
* <li>{@link AllPassFilter}</li>
2527
* </ul>
2628
*/
2729
public class SumologicMessageModelPipeline implements
28-
IKinesisConnectorPipeline<KinesisMessageModel, String> {
30+
IKinesisConnectorPipeline<SimpleKinesisMessageModel, String> {
2931

32+
private static final Log LOG = LogFactory.getLog(SumologicMessageModelPipeline.class);
33+
3034
@Override
3135
public IEmitter<String> getEmitter(KinesisConnectorConfiguration configuration) {
3236
return new SumologicEmitter(configuration);
3337
}
3438

3539
@Override
36-
public IBuffer<KinesisMessageModel> getBuffer(KinesisConnectorConfiguration configuration) {
37-
return new BasicMemoryBuffer<KinesisMessageModel>(configuration);
40+
public IBuffer<SimpleKinesisMessageModel> getBuffer(KinesisConnectorConfiguration configuration) {
41+
return new BasicMemoryBuffer<SimpleKinesisMessageModel>(configuration);
3842
}
3943

4044
@Override
41-
public ITransformer<KinesisMessageModel, String>
45+
public ITransformer<SimpleKinesisMessageModel, String>
4246
getTransformer(KinesisConnectorConfiguration configuration) {
43-
return new KinesisMessageModelSumologicTransformer();
47+
48+
// Load specified class
49+
String argClass = ((KinesisConnectorForSumologicConfiguration)configuration).TRANSFORMER_CLASS;
50+
String className = "com.mcplusa.sumologic."+argClass;
51+
ClassLoader classLoader = SumologicMessageModelPipeline.class.getClassLoader();
52+
Class ModelClass = null;
53+
try {
54+
ModelClass = classLoader.loadClass(className);
55+
ITransformer<SimpleKinesisMessageModel, String> ITransformerObject = (ITransformer<SimpleKinesisMessageModel, String>)ModelClass.newInstance();
56+
LOG.info("Using transformer: "+ITransformerObject.getClass().getName());
57+
return ITransformerObject;
58+
} catch (ClassNotFoundException e) {
59+
LOG.error("Class not found: "+className+" error: "+e.getMessage());
60+
} catch (InstantiationException e) {
61+
LOG.error("Class not found: "+className+" error: "+e.getMessage());
62+
} catch (IllegalAccessException e) {
63+
LOG.error("Class not found: "+className+" error: "+e.getMessage());
64+
}
65+
66+
return new DefaultKinesisMessageModelSumologicTransformer();
4467
}
4568

4669
@Override
47-
public IFilter<KinesisMessageModel> getFilter(KinesisConnectorConfiguration configuration) {
48-
return new AllPassFilter<KinesisMessageModel>();
70+
public IFilter<SimpleKinesisMessageModel> getFilter(KinesisConnectorConfiguration configuration) {
71+
return new AllPassFilter<SimpleKinesisMessageModel>();
4972
}
5073

5174
}

0 commit comments

Comments
 (0)