Skip to content

Commit 6bc61b3

Browse files
Juan Pablo Diaz-VazJuan Pablo Diaz-Vaz
authored andcommitted
SUMOK-18 Added transformer into Properties file
1 parent 5e76d15 commit 6bc61b3

7 files changed

Lines changed: 110 additions & 79 deletions

File tree

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;
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.KinesisMessageModel;
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 KinesisMessageModel} records in JSON format. The output is in a format
26+
* usable for insertions to Sumologic.
27+
*/
28+
public class DefaultKinesisMessageModelSumologicTransformer implements
29+
SumologicTransformer<KinesisMessageModel> {
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(KinesisMessageModel message) {
42+
return message.toString();
43+
}
44+
45+
@Override
46+
public KinesisMessageModel toClass(Record record) throws IOException {
47+
byte[] decodedRecord = record.getData().array();
48+
String stringifiedRecord = new String(decodedRecord);
49+
50+
return new KinesisMessageModel(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/SumologicMessageModelPipeline.java

Lines changed: 25 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,11 @@
11
package com.mcplusa.sumologic;
22

3+
import org.apache.commons.logging.Log;
4+
import org.apache.commons.logging.LogFactory;
5+
36
import com.mcplusa.sumologic.KinesisMessageModel;
47
import com.mcplusa.sumologic.KinesisMessageModelSumologicTransformer;
58
import com.mcplusa.sumologic.implementations.SumologicEmitter;
6-
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;
@@ -27,6 +29,8 @@
2729
public class SumologicMessageModelPipeline implements
2830
IKinesisConnectorPipeline<KinesisMessageModel, 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);
@@ -40,7 +44,26 @@ public IBuffer<KinesisMessageModel> getBuffer(KinesisConnectorConfiguration conf
4044
@Override
4145
public ITransformer<KinesisMessageModel, 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<KinesisMessageModel, String> ITransformerObject = (ITransformer<KinesisMessageModel, 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

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

Lines changed: 15 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,10 @@
1515
import org.apache.http.util.EntityUtils;
1616
import org.apache.http.impl.client.DefaultHttpClient;
1717
import org.apache.http.impl.conn.tsccm.ThreadSafeClientConnManager;
18-
import org.apache.log4j.Logger;
19-
import org.apache.log4j.PatternLayout;
20-
import org.apache.log4j.RollingFileAppender;
2118

2219
public class SumologicSender {
2320
private static final Log LOG = LogFactory.getLog(SumologicSender.class);
24-
21+
2522
private String url = null;
2623
private HttpClient httpClient = null;
2724

@@ -30,33 +27,23 @@ public class SumologicSender {
3027
private static final int RETRIES = 3;
3128
private static final int SLEEP_TIME = 1000;
3229

33-
private boolean useLog4j = false;
34-
35-
public SumologicSender(String url, boolean useLog4j) {
30+
31+
public SumologicSender(String url) {
3632
this.url = url;
37-
this.useLog4j = useLog4j;
3833

39-
if (!useLog4j) {
40-
HttpParams params = new BasicHttpParams();
41-
HttpConnectionParams.setConnectionTimeout(params, connectionTimeout);
42-
HttpConnectionParams.setSoTimeout(params, socketTimeout);
43-
httpClient = new DefaultHttpClient(new ThreadSafeClientConnManager(), params);
44-
}
45-
}
46-
47-
public boolean sendToSumologicUsingLog4j(String data) {
48-
Logger sumologicLog = Logger.getLogger("sumologic");
49-
sumologicLog.trace(data);
50-
return true;
34+
HttpParams params = new BasicHttpParams();
35+
HttpConnectionParams.setConnectionTimeout(params, connectionTimeout);
36+
HttpConnectionParams.setSoTimeout(params, socketTimeout);
37+
httpClient = new DefaultHttpClient(new ThreadSafeClientConnManager(), params);
5138
}
52-
53-
public boolean sendToSumologicUsingHTTPRequest(String data) throws IOException {
39+
40+
public boolean sendToSumologic(String data) throws IOException{
5441
int retries = RETRIES;
55-
int sleep_time = SLEEP_TIME;
56-
int statusCode;
57-
58-
do {
59-
HttpPost post = null;
42+
int sleep_time = SLEEP_TIME;
43+
int statusCode;
44+
45+
do {
46+
HttpPost post = null;
6047
post = new HttpPost(url);
6148
post.setEntity(new StringEntity(data, HTTP.PLAIN_TEXT_TYPE, HTTP.UTF_8));
6249
HttpResponse response = httpClient.execute(post);
@@ -74,7 +61,7 @@ public boolean sendToSumologicUsingHTTPRequest(String data) throws IOException {
7461
Thread.sleep(sleep_time);
7562
} catch (InterruptedException ignore) {}
7663
}
77-
} while (statusCode == 429 && retries > 0);
64+
} while (statusCode == 429 && retries > 0);
7865

7966
// Check if the request was successful;
8067
if (statusCode != 200) {
@@ -85,13 +72,4 @@ public boolean sendToSumologicUsingHTTPRequest(String data) throws IOException {
8572
return true;
8673
}
8774
}
88-
89-
public boolean sendToSumologic(String data) throws IOException{
90-
if (this.useLog4j) {
91-
return sendToSumologicUsingLog4j(data);
92-
}
93-
else {
94-
return sendToSumologicUsingHTTPRequest(data);
95-
}
96-
}
9775
}

src/main/java/com/mcplusa/sumologic/implementations/SumologicEmitter.java

Lines changed: 4 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -23,29 +23,24 @@
2323
*/
2424
public class SumologicEmitter implements IEmitter<String> {
2525
private static final Log LOG = LogFactory.getLog(SumologicEmitter.class);
26-
27-
private static final boolean CONCATENATE_BATCH = false;
2826

2927
private SumologicSender sender;
3028
private KinesisConnectorForSumologicConfiguration config;
3129

3230
public SumologicEmitter(KinesisConnectorConfiguration configuration) {
3331
this.config = (KinesisConnectorForSumologicConfiguration) configuration;
34-
sender = new SumologicSender(config.SUMOLOGIC_URL, config.SUMOLOGIC_USE_LOG4J);
32+
sender = new SumologicSender(this.config.SUMOLOGIC_URL);
3533
}
3634

37-
public SumologicEmitter(String url, boolean useLog4j) {
38-
sender = new SumologicSender(url, useLog4j);
35+
public SumologicEmitter(String url) {
36+
sender = new SumologicSender(url);
3937
}
4038

4139
@Override
4240
public List<String> emit(final UnmodifiableBuffer<String> buffer)
4341
throws IOException {
4442
List<String> records = buffer.getRecords();
45-
if (CONCATENATE_BATCH)
46-
return sendBatchConcatenating(records);
47-
else
48-
return sendRecordsOneByOne(records);
43+
return sendBatchConcatenating(records);
4944
}
5045

5146
public List<String> sendBatchConcatenating(List<String> records) {
@@ -70,21 +65,6 @@ public List<String> sendBatchConcatenating(List<String> records) {
7065
return records;
7166
}
7267
}
73-
74-
public List<String> sendRecordsOneByOne (List<String> records) {
75-
ArrayList<String> failedRecords = new ArrayList<String>();
76-
for (String record: records) {
77-
try {
78-
if (!sender.sendToSumologic(record)) {
79-
failedRecords.add(record);
80-
}
81-
} catch (IOException e) {
82-
LOG.warn("Couldn't send record: "+record);
83-
}
84-
}
85-
LOG.info("Sent records: "+(records.size()-failedRecords.size())+" failed: "+failedRecords.size());
86-
return failedRecords;
87-
}
8868

8969
@Override
9070
public void fail(List<String> records) {

src/test/java/com/mcplusa/sumologic/implementations/SumologicEmitterTest.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ public void theEmitterShouldReturnTheListParameterWhenFailing () {
3939
messages.add("This is message #3");
4040
messages.add("This is message #4");
4141

42-
SumologicEmitter emitter = new SumologicEmitter(url, false);
42+
SumologicEmitter emitter = new SumologicEmitter(url);
4343
List <String> notEmittedMessages = emitter.sendBatchConcatenating(messages);
4444

4545
Assert.assertEquals(messages, notEmittedMessages);
@@ -55,7 +55,7 @@ public void theEmitterShouldReturnAnEmptyListOnSuccess () {
5555
messages.add("This is message #3");
5656
messages.add("This is message #4");
5757

58-
SumologicEmitter emitter = new SumologicEmitter(url, false);
58+
SumologicEmitter emitter = new SumologicEmitter(url);
5959
List <String> notEmittedMessages = emitter.sendBatchConcatenating(messages);
6060

6161
Assert.assertEquals(0, notEmittedMessages.size());

0 commit comments

Comments
 (0)