Skip to content

Commit 5e76d15

Browse files
committed
Merged in SUMOK-17 (pull request #5)
Sumok 17
2 parents 70cabd4 + 8bd1b13 commit 5e76d15

5 files changed

Lines changed: 90 additions & 22 deletions

File tree

log4j.properties

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
# Root logger option
2+
log4j.rootLogger=INFO, stdout
3+
log4j.logger.sumologic = TRACE, sumo
4+
5+
# Direct log messages to sumo
6+
log4j.appender.sumo=com.sumologic.log4j.BufferedSumoLogicAppender
7+
log4j.appender.sumo.url=https://collectors.us2.sumologic.com/receiver/v1/http/ZaVnC4dhaV0GzIY4tZaKLL26afV52gXBvSFc3jG1eLc2lKINzS2doZdRjUMQMb2CXK8r6fdmHoUazJiHjJ-2OygApoWFaxCTkWFrzAiraCc5i411pkio-g==
8+
log4j.appender.sumo.layout=org.apache.log4j.PatternLayout
9+
log4j.appender.sumo.layout.ConversionPattern=%d{DATE} %5p %c{1}:%L - %m%n
10+
log4j.additivity.sumo = false
11+
log4j.appender.sumo.Threshold = TRACE
12+
13+
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
14+
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
15+
log4j.appender.stdout.layout.ConversionPattern=%d{DATE} %5p %c{1}:%L - %m%n
16+
log4j.appender.stdout.Threshold = INFO

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

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,15 +14,25 @@
1414
public class KinesisConnectorForSumologicConfiguration extends KinesisConnectorConfiguration {
1515
// Properties added for Sumologic
1616
public static final String PROP_SUMOLOGIC_URL = "sumologicUrl";
17+
public static final String PROP_SUMOLOGIC_USE_LOG4J = "useLog4j";
1718
public final String SUMOLOGIC_URL;
19+
public final boolean SUMOLOGIC_USE_LOG4J;
20+
21+
public final boolean DEFAULT_SUMOLOGIC_USE_LOG4J = false;
1822

1923
/**
2024
* Configure the connector application with any set of properties that are unique to the application. Any
2125
* unspecified property will be set to a default value.
2226
*/
2327
public KinesisConnectorForSumologicConfiguration(Properties properties, AWSCredentialsProvider credentialsProvider) {
2428
super(properties, credentialsProvider);
25-
2629
SUMOLOGIC_URL = properties.getProperty(PROP_SUMOLOGIC_URL, null);
30+
SUMOLOGIC_USE_LOG4J = getBooleanProperty(PROP_SUMOLOGIC_USE_LOG4J,
31+
DEFAULT_SUMOLOGIC_USE_LOG4J, properties);
2732
}
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+
}
2838
}

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

Lines changed: 37 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -15,10 +15,13 @@
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;
1821

1922
public class SumologicSender {
2023
private static final Log LOG = LogFactory.getLog(SumologicSender.class);
21-
24+
2225
private String url = null;
2326
private HttpClient httpClient = null;
2427

@@ -27,23 +30,33 @@ public class SumologicSender {
2730
private static final int RETRIES = 3;
2831
private static final int SLEEP_TIME = 1000;
2932

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

34-
HttpParams params = new BasicHttpParams();
35-
HttpConnectionParams.setConnectionTimeout(params, connectionTimeout);
36-
HttpConnectionParams.setSoTimeout(params, socketTimeout);
37-
httpClient = new DefaultHttpClient(new ThreadSafeClientConnManager(), params);
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+
}
3845
}
39-
40-
public boolean sendToSumologic(String data) throws IOException{
46+
47+
public boolean sendToSumologicUsingLog4j(String data) {
48+
Logger sumologicLog = Logger.getLogger("sumologic");
49+
sumologicLog.trace(data);
50+
return true;
51+
}
52+
53+
public boolean sendToSumologicUsingHTTPRequest(String data) throws IOException {
4154
int retries = RETRIES;
42-
int sleep_time = SLEEP_TIME;
43-
int statusCode;
44-
45-
do {
46-
HttpPost post = null;
55+
int sleep_time = SLEEP_TIME;
56+
int statusCode;
57+
58+
do {
59+
HttpPost post = null;
4760
post = new HttpPost(url);
4861
post.setEntity(new StringEntity(data, HTTP.PLAIN_TEXT_TYPE, HTTP.UTF_8));
4962
HttpResponse response = httpClient.execute(post);
@@ -61,7 +74,7 @@ public boolean sendToSumologic(String data) throws IOException{
6174
Thread.sleep(sleep_time);
6275
} catch (InterruptedException ignore) {}
6376
}
64-
} while (statusCode == 429 && retries > 0);
77+
} while (statusCode == 429 && retries > 0);
6578

6679
// Check if the request was successful;
6780
if (statusCode != 200) {
@@ -72,4 +85,13 @@ public boolean sendToSumologic(String data) throws IOException{
7285
return true;
7386
}
7487
}
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+
}
7597
}

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

Lines changed: 24 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,24 +23,29 @@
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;
2628

2729
private SumologicSender sender;
2830
private KinesisConnectorForSumologicConfiguration config;
2931

3032
public SumologicEmitter(KinesisConnectorConfiguration configuration) {
3133
this.config = (KinesisConnectorForSumologicConfiguration) configuration;
32-
sender = new SumologicSender(this.config.SUMOLOGIC_URL);
34+
sender = new SumologicSender(config.SUMOLOGIC_URL, config.SUMOLOGIC_USE_LOG4J);
3335
}
3436

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

3941
@Override
4042
public List<String> emit(final UnmodifiableBuffer<String> buffer)
4143
throws IOException {
4244
List<String> records = buffer.getRecords();
43-
return sendBatchConcatenating(records);
45+
if (CONCATENATE_BATCH)
46+
return sendBatchConcatenating(records);
47+
else
48+
return sendRecordsOneByOne(records);
4449
}
4550

4651
public List<String> sendBatchConcatenating(List<String> records) {
@@ -65,6 +70,21 @@ public List<String> sendBatchConcatenating(List<String> records) {
6570
return records;
6671
}
6772
}
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+
}
6888

6989
@Override
7090
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);
42+
SumologicEmitter emitter = new SumologicEmitter(url, false);
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);
58+
SumologicEmitter emitter = new SumologicEmitter(url, false);
5959
List <String> notEmittedMessages = emitter.sendBatchConcatenating(messages);
6060

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

0 commit comments

Comments
 (0)