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 ;
30+ import com .google .gson .Gson ;
31+
2032/**
21- * 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
2234 * usable for insertions to Sumologic.
2335 */
24- public class CloudWatchMessageModelSumologicTransformer implements
25- SumologicTransformer <SimpleKinesisMessageModel > {
36+ public class CloudWatchMessageModelSumologicTransformer
37+ implements SumologicTransformer <CloudWatchLogsMessageModel > {
2638
2739 private static final Log LOG = LogFactory .getLog (CloudWatchMessageModelSumologicTransformer .class );
40+
41+ private static CharsetEncoder encoder = Charset .forName ("UTF-8" ).newEncoder ();
2842
2943 /**
3044 * Creates a new KinesisMessageModelSumologicTransformer.
@@ -34,24 +48,70 @@ public CloudWatchMessageModelSumologicTransformer() {
3448 }
3549
3650 @ Override
37- public String fromClass (SimpleKinesisMessageModel message ) {
38- return message .toString ();
39- }
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 ()));
4068
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+
4184 @ Override
42- public SimpleKinesisMessageModel toClass (Record record ) throws IOException {
85+ public CloudWatchLogsMessageModel toClass (Record record ) {
4386 byte [] decodedRecord = record .getData ().array ();
4487 String stringifiedRecord = decompressGzip (decodedRecord );
4588
4689 if (stringifiedRecord == null ) {
47- LOG .error ("Unable to decompress the record: " +new String (record .getData ().array ()));
48- LOG . error ( "Not attempting to transform into a Message Model" );
90+ LOG .error ("Unable to decompress the record: " +new String (record .getData ().array ())
91+ + " \n Not attempting to transform into a Message Model" );
4992 return null ;
5093 }
94+
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 ());
100+ }
101+ record .setData (bufferedData );
51102
52- 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+
53113 }
54-
114+
55115 public static String decompressGzip (byte [] compressedData ) {
56116 try {
57117 GZIPInputStream gis = new GZIPInputStream (new ByteArrayInputStream (compressedData ));
@@ -83,4 +143,14 @@ public static String byteBufferToString(ByteBuffer buffer){
83143 return data ;
84144 }
85145
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+
86156}
0 commit comments