Skip to content

Commit 86915ac

Browse files
Juan Pablo Diaz-VazJuan Pablo Diaz-Vaz
authored andcommitted
Separated models into separate package and utils into new class
1 parent af18b58 commit 86915ac

14 files changed

Lines changed: 163 additions & 103 deletions

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

Lines changed: 3 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,9 @@
77
import com.amazonaws.util.json.JSONArray;
88
import com.amazonaws.util.json.JSONException;
99
import com.amazonaws.util.json.JSONObject;
10-
import com.sumologic.client.SimpleKinesisMessageModel;
1110
import com.sumologic.client.implementations.SumologicTransformer;
11+
import com.sumologic.client.model.CloudWatchLogsMessageModel;
12+
import com.sumologic.client.model.LogEvent;
1213

1314
import org.apache.commons.logging.Log;
1415
import org.apache.commons.logging.LogFactory;
@@ -84,7 +85,7 @@ public String fromClass(CloudWatchLogsMessageModel message) {
8485
@Override
8586
public CloudWatchLogsMessageModel toClass(Record record) {
8687
byte[] decodedRecord = record.getData().array();
87-
String stringifiedRecord = decompressGzip(decodedRecord);
88+
String stringifiedRecord = SumologicKinesisUtils.decompressGzip(decodedRecord);
8889

8990
if (stringifiedRecord == null) {
9091
LOG.error("Unable to decompress the record: "+new String(record.getData().array())
@@ -109,48 +110,5 @@ public CloudWatchLogsMessageModel toClass(Record record) {
109110
+"\nerror: "+e.getMessage());
110111
}
111112
return null;
112-
113-
}
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-
}
130113
}
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-
156114
}

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,9 +4,9 @@
44

55
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
66
import com.amazonaws.services.kinesis.model.Record;
7-
import com.sumologic.client.SimpleKinesisMessageModel;
87
import com.sumologic.client.implementations.SumologicEmitter;
98
import com.sumologic.client.implementations.SumologicTransformer;
9+
import com.sumologic.client.model.SimpleKinesisMessageModel;
1010

1111
import org.apache.commons.logging.Log;
1212
import org.apache.commons.logging.LogFactory;

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

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

33
import com.amazonaws.services.kinesis.connectors.KinesisConnectorRecordProcessorFactory;
4-
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> {
109

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

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,8 @@
33
import org.apache.commons.logging.Log;
44
import org.apache.commons.logging.LogFactory;
55

6-
import com.sumologic.client.SimpleKinesisMessageModel;
7-
import com.sumologic.client.CloudWatchMessageModelSumologicTransformer;
86
import com.sumologic.client.implementations.SumologicEmitter;
7+
import com.sumologic.client.model.SimpleKinesisMessageModel;
98
import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
109
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
1110
import com.amazonaws.services.kinesis.connectors.impl.BasicMemoryBuffer;

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/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;

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

Lines changed: 1 addition & 1 deletion
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.HashMap;
44
import java.util.Map;

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

Lines changed: 1 addition & 1 deletion
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.io.Serializable;
44

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,8 +14,7 @@
1414
import org.apache.commons.logging.Log;
1515
import org.apache.commons.logging.LogFactory;
1616

17-
import com.sumologic.client.SimpleKinesisMessageModel;
18-
17+
import com.sumologic.client.model.SimpleKinesisMessageModel;
1918
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
2019
import com.amazonaws.services.kinesis.model.PutRecordRequest;
2120

0 commit comments

Comments
 (0)