Skip to content

Commit c5119d1

Browse files
committed
Merge pull request #3 from jpdiazvaz/master
Better handling of errors during execution, using Log4j for logging
2 parents 71939cf + 1726448 commit c5119d1

21 files changed

Lines changed: 438 additions & 172 deletions

build.xml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@
3434
<get src= "http://central.maven.org/maven2/commons-codec/commons-codec/1.10/commons-codec-1.10.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
3535
<get src= "http://central.maven.org/maven2/com/ning/async-http-client/1.9.30/async-http-client-1.9.30.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
3636
<get src= "http://central.maven.org/maven2/com/google/code/gson/gson/2.3.1/gson-2.3.1.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
37+
<get src= "http://central.maven.org/maven2/log4j/log4j/1.2.17/log4j-1.2.17.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
3738

3839
<get src= "https://hamcrest.googlecode.com/files/hamcrest-core-1.3.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
3940

log4j.properties

Lines changed: 0 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,6 @@
11
# Root logger option
22
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
123

134
log4j.appender.stdout=org.apache.log4j.ConsoleAppender
145
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
156
log4j.appender.stdout.layout.ConversionPattern=%d{DATE} %5p %c{1}:%L - %m%n
16-
log4j.appender.stdout.Threshold = INFO

src/main/java/com/sumologic/client/CloudWatchMessageModelSumologicTransformer.java

Lines changed: 6 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -2,42 +2,33 @@
22

33
import java.io.IOException;
44

5-
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
65
import com.amazonaws.services.kinesis.model.Record;
76
import com.amazonaws.util.json.JSONArray;
87
import com.amazonaws.util.json.JSONException;
98
import com.amazonaws.util.json.JSONObject;
10-
import com.sumologic.client.SimpleKinesisMessageModel;
119
import com.sumologic.client.implementations.SumologicTransformer;
10+
import com.sumologic.client.model.CloudWatchLogsMessageModel;
11+
import com.sumologic.client.model.LogEvent;
1212

13-
import org.apache.commons.logging.Log;
14-
import org.apache.commons.logging.LogFactory;
13+
import org.apache.log4j.Logger;
1514

16-
import java.io.ByteArrayInputStream;
17-
import java.io.BufferedReader;
18-
import java.io.InputStreamReader;
1915
import java.nio.ByteBuffer;
2016
import java.nio.CharBuffer;
2117
import java.nio.charset.CharacterCodingException;
2218
import java.nio.charset.Charset;
23-
import java.nio.charset.CharsetDecoder;
2419
import java.nio.charset.CharsetEncoder;
2520
import java.util.List;
26-
import java.util.zip.GZIPInputStream;
2721

28-
import com.fasterxml.jackson.databind.JsonMappingException;
2922
import com.fasterxml.jackson.databind.ObjectMapper;
30-
import com.google.gson.Gson;
3123

3224
/**
3325
* A custom transfomer for {@link CloudWatchLogsMessageModel} records in JSON format. The output is in a format
3426
* usable for insertions to Sumologic.
3527
*/
3628
public class CloudWatchMessageModelSumologicTransformer
3729
implements SumologicTransformer<CloudWatchLogsMessageModel> {
38-
39-
private static final Log LOG = LogFactory.getLog(CloudWatchMessageModelSumologicTransformer.class);
40-
30+
private static final Logger LOG = Logger.getLogger(CloudWatchMessageModelSumologicTransformer.class.getName());
31+
4132
private static CharsetEncoder encoder = Charset.forName("UTF-8").newEncoder();
4233

4334
/**
@@ -84,7 +75,7 @@ public String fromClass(CloudWatchLogsMessageModel message) {
8475
@Override
8576
public CloudWatchLogsMessageModel toClass(Record record) {
8677
byte[] decodedRecord = record.getData().array();
87-
String stringifiedRecord = decompressGzip(decodedRecord);
78+
String stringifiedRecord = SumologicKinesisUtils.decompressGzip(decodedRecord);
8879

8980
if (stringifiedRecord == null) {
9081
LOG.error("Unable to decompress the record: "+new String(record.getData().array())
@@ -109,48 +100,5 @@ public CloudWatchLogsMessageModel toClass(Record record) {
109100
+"\nerror: "+e.getMessage());
110101
}
111102
return null;
112-
113103
}
114-
115-
public static String decompressGzip(byte[] compressedData) {
116-
try {
117-
GZIPInputStream gis = new GZIPInputStream(new ByteArrayInputStream(compressedData));
118-
BufferedReader bf = new BufferedReader(new InputStreamReader(gis, "UTF-8"));
119-
120-
String outStr = "";
121-
String line;
122-
while ((line=bf.readLine())!=null) {
123-
outStr += line;
124-
}
125-
return outStr;
126-
} catch (IOException exc) {
127-
LOG.warn("Exception during decompression of data: " + exc.getMessage());
128-
return null;
129-
}
130-
}
131-
132-
public static String byteBufferToString(ByteBuffer buffer){
133-
String data = "";
134-
CharsetDecoder decoder = Charset.forName("UTF-8").newDecoder();
135-
try{
136-
int old_position = buffer.position();
137-
data = decoder.decode(buffer).toString();
138-
buffer.position(old_position);
139-
}catch (Exception e){
140-
e.printStackTrace();
141-
return "";
142-
}
143-
return data;
144-
}
145-
146-
private static final Gson gson = new Gson();
147-
public static boolean verifyJSON(String json) {
148-
try {
149-
gson.fromJson(json, Object.class);
150-
return true;
151-
} catch(com.google.gson.JsonSyntaxException ex) {
152-
return false;
153-
}
154-
}
155-
156104
}

src/main/java/com/sumologic/client/DefaultKinesisMessageModelSumologicTransformer.java

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

33
import java.io.IOException;
44

5-
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
65
import com.amazonaws.services.kinesis.model.Record;
7-
import com.sumologic.client.SimpleKinesisMessageModel;
8-
import com.sumologic.client.implementations.SumologicEmitter;
96
import com.sumologic.client.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-
7+
import com.sumologic.client.model.SimpleKinesisMessageModel;
238

249
/**
2510
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
2611
* usable for insertions to Sumologic.
2712
*/
2813
public class DefaultKinesisMessageModelSumologicTransformer implements
2914
SumologicTransformer<SimpleKinesisMessageModel> {
30-
31-
private static final Log LOG = LogFactory.getLog(DefaultKinesisMessageModelSumologicTransformer.class);
32-
3315
/**
3416
* Creates a new KinesisMessageModelSumologicTransformer.
3517
*/
@@ -46,7 +28,7 @@ public String fromClass(SimpleKinesisMessageModel message) {
4628
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
4729
byte[] decodedRecord = record.getData().array();
4830
String stringifiedRecord = new String(decodedRecord);
49-
31+
5032
return new SimpleKinesisMessageModel(stringifiedRecord);
5133
}
5234
}

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

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

3-
import com.amazonaws.services.kinesis.connectors.KinesisConnectorRecordProcessorFactory;
4-
3+
import com.sumologic.kinesis.KinesisConnectorRecordProcessorFactory;
54
import com.sumologic.kinesis.KinesisConnectorExecutor;
6-
import com.sumologic.client.SimpleKinesisMessageModel;
75
import com.sumologic.client.SumologicMessageModelPipeline;
6+
import com.sumologic.client.model.SimpleKinesisMessageModel;
87

98
public class SumologicExecutor extends KinesisConnectorExecutor<SimpleKinesisMessageModel, String> {
10-
119
private static String configFile = "SumologicConnector.properties";
1210

1311
/**
Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
package com.sumologic.client;
2+
3+
import java.io.BufferedReader;
4+
import java.io.ByteArrayInputStream;
5+
import java.io.ByteArrayOutputStream;
6+
import java.io.IOException;
7+
import java.io.InputStreamReader;
8+
import java.nio.ByteBuffer;
9+
import java.nio.charset.Charset;
10+
import java.nio.charset.CharsetDecoder;
11+
import java.util.zip.GZIPInputStream;
12+
import java.util.zip.GZIPOutputStream;
13+
14+
import org.apache.log4j.Logger;
15+
16+
import com.google.gson.Gson;
17+
18+
public class SumologicKinesisUtils {
19+
private static final Logger LOG = Logger.getLogger(SumologicKinesisUtils.class.getName());
20+
21+
public static byte[] compressGzip(String data) {
22+
if (data == null || data.length() == 0) {
23+
return null;
24+
}
25+
26+
ByteArrayOutputStream outputStream=new ByteArrayOutputStream();
27+
GZIPOutputStream gzip;
28+
try {
29+
gzip = new GZIPOutputStream(outputStream);
30+
} catch (IOException e) {
31+
LOG.error("Cannot compress into GZIP "+e.getMessage());
32+
return null;
33+
}
34+
35+
// Put data into the GZIP buffer
36+
try {
37+
gzip.write(data.getBytes("UTF-8"));
38+
gzip.close();
39+
} catch (IOException e) {
40+
e.printStackTrace();
41+
}
42+
43+
return outputStream.toByteArray();
44+
}
45+
46+
public static String decompressGzip(byte[] compressedData) {
47+
try {
48+
GZIPInputStream gis = new GZIPInputStream(new ByteArrayInputStream(compressedData));
49+
BufferedReader bf = new BufferedReader(new InputStreamReader(gis, "UTF-8"));
50+
51+
String outStr = "";
52+
String line;
53+
while ((line=bf.readLine())!=null) {
54+
outStr += line;
55+
}
56+
return outStr;
57+
} catch (IOException exc) {
58+
LOG.warn("Exception during decompression of data: " + exc.getMessage());
59+
return null;
60+
}
61+
}
62+
63+
public static String byteBufferToString(ByteBuffer buffer){
64+
String data = "";
65+
CharsetDecoder decoder = Charset.forName("UTF-8").newDecoder();
66+
try{
67+
int old_position = buffer.position();
68+
data = decoder.decode(buffer).toString();
69+
buffer.position(old_position);
70+
}catch (Exception e){
71+
e.printStackTrace();
72+
return "";
73+
}
74+
return data;
75+
}
76+
77+
private static final Gson gson = new Gson();
78+
public static boolean verifyJSON(String json) {
79+
try {
80+
gson.fromJson(json, Object.class);
81+
return true;
82+
} catch(com.google.gson.JsonSyntaxException ex) {
83+
return false;
84+
}
85+
}
86+
}

src/main/java/com/sumologic/client/SumologicMessageModelPipeline.java

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

3-
import org.apache.commons.logging.Log;
4-
import org.apache.commons.logging.LogFactory;
3+
import org.apache.log4j.Logger;
54

6-
import com.sumologic.client.SimpleKinesisMessageModel;
7-
import com.sumologic.client.CloudWatchMessageModelSumologicTransformer;
85
import com.sumologic.client.implementations.SumologicEmitter;
6+
import com.sumologic.client.model.SimpleKinesisMessageModel;
97
import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
108
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
119
import com.amazonaws.services.kinesis.connectors.impl.BasicMemoryBuffer;
@@ -28,9 +26,8 @@
2826
*/
2927
public class SumologicMessageModelPipeline implements
3028
IKinesisConnectorPipeline<SimpleKinesisMessageModel, String> {
29+
private static final Logger LOG = Logger.getLogger(SumologicMessageModelPipeline.class.getName());
3130

32-
private static final Log LOG = LogFactory.getLog(SumologicMessageModelPipeline.class);
33-
3431
@Override
3532
public IEmitter<String> getEmitter(KinesisConnectorConfiguration configuration) {
3633
return new SumologicEmitter(configuration);

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

Lines changed: 1 addition & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ public boolean sendToSumologic(String data) throws IOException{
4848
BoundRequestBuilder builder = null;
4949
builder = this.clientPreparePost(url);
5050

51-
byte[] compressedData = SumologicSender.compressGzip(data);
51+
byte[] compressedData = SumologicKinesisUtils.compressGzip(data);
5252

5353
builder.setHeader("Content-Encoding", "gzip");
5454
builder.setBody(compressedData);
@@ -81,29 +81,4 @@ public boolean sendToSumologic(String data) throws IOException{
8181
return true;
8282
}
8383
}
84-
85-
public static byte[] compressGzip(String data) {
86-
if (data == null || data.length() == 0) {
87-
return null;
88-
}
89-
90-
ByteArrayOutputStream outputStream=new ByteArrayOutputStream();
91-
GZIPOutputStream gzip;
92-
try {
93-
gzip = new GZIPOutputStream(outputStream);
94-
} catch (IOException e) {
95-
LOG.error("Cannot compress into GZIP "+e.getMessage());
96-
return null;
97-
}
98-
99-
// Put data into the GZIP buffer
100-
try {
101-
gzip.write(data.getBytes("UTF-8"));
102-
gzip.close();
103-
} catch (IOException e) {
104-
e.printStackTrace();
105-
}
106-
107-
return outputStream.toByteArray();
108-
}
10984
}

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

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

33
import java.io.IOException;
44
import java.util.ArrayList;
5-
import java.util.HashSet;
65
import java.util.List;
7-
import java.util.Set;
86

9-
import org.apache.commons.logging.Log;
10-
import org.apache.commons.logging.LogFactory;
7+
import org.apache.log4j.Logger;
118

129
import com.sumologic.client.SumologicSender;
1310
import com.sumologic.client.KinesisConnectorForSumologicConfiguration;
14-
1511
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
1612
import com.amazonaws.services.kinesis.connectors.UnmodifiableBuffer;
1713
import com.amazonaws.services.kinesis.connectors.interfaces.IEmitter;
@@ -22,7 +18,7 @@
2218
* Sumologic.
2319
*/
2420
public class SumologicEmitter implements IEmitter<String> {
25-
private static final Log LOG = LogFactory.getLog(SumologicEmitter.class);
21+
private static final Logger LOG = Logger.getLogger(SumologicEmitter.class.getName());
2622

2723
private SumologicSender sender;
2824
private KinesisConnectorForSumologicConfiguration config;

src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java renamed to src/main/java/com/sumologic/client/model/CloudWatchLogsMessageModel.java

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
package com.sumologic.client;
1+
package com.sumologic.client.model;
22

33
import java.util.ArrayList;
44
import java.util.HashMap;
@@ -7,7 +7,6 @@
77

88
import org.apache.commons.lang.builder.ToStringBuilder;
99

10-
import com.sumologic.client.LogEvent;
1110
import com.fasterxml.jackson.annotation.JsonAnyGetter;
1211
import com.fasterxml.jackson.annotation.JsonAnySetter;
1312
import com.fasterxml.jackson.annotation.JsonIgnore;

0 commit comments

Comments
 (0)