22
33import java .io .IOException ;
44
5+ import com .amazonaws .services .kinesis .connectors .BasicJsonTransformer ;
56import com .amazonaws .services .kinesis .model .Record ;
7+ import com .amazonaws .util .json .JSONArray ;
8+ import com .amazonaws .util .json .JSONException ;
9+ import com .amazonaws .util .json .JSONObject ;
610import com .sumologic .client .SimpleKinesisMessageModel ;
711import com .sumologic .client .implementations .SumologicTransformer ;
812
1317import java .io .BufferedReader ;
1418import java .io .InputStreamReader ;
1519import java .nio .ByteBuffer ;
20+ import java .nio .CharBuffer ;
21+ import java .nio .charset .CharacterCodingException ;
1622import java .nio .charset .Charset ;
1723import java .nio .charset .CharsetDecoder ;
24+ import java .nio .charset .CharsetEncoder ;
25+ import java .util .List ;
1826import java .util .zip .GZIPInputStream ;
1927
28+ import com .fasterxml .jackson .databind .JsonMappingException ;
29+ import com .fasterxml .jackson .databind .ObjectMapper ;
2030import com .google .gson .Gson ;
2131
2232/**
23- * A custom transfomer for {@link SimpleKinesisMessageModel } records in JSON format. The output is in a format
33+ * A custom transfomer for {@link CloudWatchLogsMessageModel } records in JSON format. The output is in a format
2434 * usable for insertions to Sumologic.
2535 */
26- public class CloudWatchMessageModelSumologicTransformer implements
27- SumologicTransformer <SimpleKinesisMessageModel > {
36+ public class CloudWatchMessageModelSumologicTransformer
37+ implements SumologicTransformer <CloudWatchLogsMessageModel > {
2838
2939 private static final Log LOG = LogFactory .getLog (CloudWatchMessageModelSumologicTransformer .class );
40+
41+ private static CharsetEncoder encoder = Charset .forName ("UTF-8" ).newEncoder ();
3042
3143 /**
3244 * Creates a new KinesisMessageModelSumologicTransformer.
@@ -36,12 +48,41 @@ public CloudWatchMessageModelSumologicTransformer() {
3648 }
3749
3850 @ Override
39- public String fromClass (SimpleKinesisMessageModel message ) {
40- return message .toString ();
41- }
51+ public String fromClass (CloudWatchLogsMessageModel message ) {
52+ String jsonMessage = "" ;
53+ JSONObject outputObject ;
54+
55+ List <LogEvent > logEvents = message .getLogEvents ();
56+ int logEventsSize = logEvents .size ();
57+ for (int i =0 ;i <logEventsSize ;i ++) {
58+ LogEvent log = logEvents .get (i );
59+
60+ outputObject = new JSONObject ();
61+ try {
62+ // Header
63+ outputObject .put ("logGroup" , message .getLogGroup ());
64+ outputObject .put ("logStream" , message .getLogStream ());
65+ outputObject .put ("messageType" , message .getMessageType ());
66+ outputObject .put ("owner" , message .getOwner ());
67+ outputObject .put ("subscriptionFilters" , new JSONArray (message .getSubscriptionFilters ()));
4268
69+ // Body
70+ outputObject .put ("id" , log .getId ());
71+ outputObject .put ("message" , log .getMessage ());
72+ outputObject .put ("timestamp" , log .getTimestamp ());
73+ } catch (JSONException e ) {
74+ LOG .error ("Unable to convert message into JSON String: " +e .getMessage ());
75+ }
76+ jsonMessage += outputObject .toString ();
77+ if (i < logEventsSize - 1 ) {
78+ jsonMessage += '\n' ;
79+ }
80+ }
81+ return jsonMessage ;
82+ }
83+
4384 @ Override
44- public SimpleKinesisMessageModel toClass (Record record ) throws IOException {
85+ public CloudWatchLogsMessageModel toClass (Record record ) {
4586 byte [] decodedRecord = record .getData ().array ();
4687 String stringifiedRecord = decompressGzip (decodedRecord );
4788
@@ -51,15 +92,26 @@ public SimpleKinesisMessageModel toClass(Record record) throws IOException {
5192 return null ;
5293 }
5394
54- if (!verifyJSON (stringifiedRecord )) {
55- LOG .error ("The record is not a valid JSON: " +stringifiedRecord
56- +"\n Not attempting to transform into a Message Model" );
57- return null ;
95+ ByteBuffer bufferedData = null ;
96+ try {
97+ bufferedData = encoder .encode (CharBuffer .wrap (stringifiedRecord ));
98+ } catch (CharacterCodingException e ) {
99+ LOG .error ("Unable to set the decompressed Record for serializing " +e .getMessage ());
58100 }
101+ record .setData (bufferedData );
59102
60- return new SimpleKinesisMessageModel (stringifiedRecord );
103+ try {
104+ return new ObjectMapper ().readValue (
105+ record .getData ().array (),
106+ CloudWatchLogsMessageModel .class );
107+ } catch (IOException e ) {
108+ LOG .error ("Unable to convert the Record into a POJO: " +stringifiedRecord
109+ +"\n error: " +e .getMessage ());
110+ }
111+ return null ;
112+
61113 }
62-
114+
63115 public static String decompressGzip (byte [] compressedData ) {
64116 try {
65117 GZIPInputStream gis = new GZIPInputStream (new ByteArrayInputStream (compressedData ));
0 commit comments