Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions build.xml
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
<get src= "http://central.maven.org/maven2/org/apache/lucene/lucene-core/4.8.1/lucene-core-4.8.1.jar" dest="${external.dir}/lib" usetimestamp="true" verbose="true"/>
<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"/>
<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"/>
<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"/>

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

Expand Down
119 changes: 119 additions & 0 deletions src/main/java/com/sumologic/client/CloudWatchLogsMessageModel.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
package com.sumologic.client;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

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

import com.sumologic.client.LogEvent;
import com.fasterxml.jackson.annotation.JsonAnyGetter;
import com.fasterxml.jackson.annotation.JsonAnySetter;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonPropertyOrder;

@JsonInclude(JsonInclude.Include.NON_NULL)
@JsonPropertyOrder({
"logEvents",
"logGroup",
"logStream",
"messageType",
"owner",
"subscriptionFilters"
})

public class CloudWatchLogsMessageModel {

@JsonProperty("logEvents")
private List<LogEvent> logEvents = new ArrayList<LogEvent>();
@JsonProperty("logGroup")
private String logGroup;
@JsonProperty("logStream")
private String logStream;
@JsonProperty("messageType")
private String messageType;
@JsonProperty("owner")
private String owner;
@JsonProperty("subscriptionFilters")
private List<String> subscriptionFilters = new ArrayList<String>();
@JsonIgnore
private Map<String, Object> additionalProperties = new HashMap<String, Object>();

@JsonProperty("logEvents")
public List<LogEvent> getLogEvents() {
return logEvents;
}

@JsonProperty("logEvents")
public void setLogEvents(List<LogEvent> logEvents) {
this.logEvents = logEvents;
}

@JsonProperty("logGroup")
public String getLogGroup() {
return logGroup;
}

@JsonProperty("logGroup")
public void setLogGroup(String logGroup) {
this.logGroup = logGroup;
}

@JsonProperty("logStream")
public String getLogStream() {
return logStream;
}

@JsonProperty("logStream")
public void setLogStream(String logStream) {
this.logStream = logStream;
}

@JsonProperty("messageType")
public String getMessageType() {
return messageType;
}

@JsonProperty("messageType")
public void setMessageType(String messageType) {
this.messageType = messageType;
}

@JsonProperty("owner")
public String getOwner() {
return owner;
}

@JsonProperty("owner")
public void setOwner(String owner) {
this.owner = owner;
}

@JsonProperty("subscriptionFilters")
public List<String> getSubscriptionFilters() {
return subscriptionFilters;
}

@JsonProperty("subscriptionFilters")
public void setSubscriptionFilters(List<String> subscriptionFilters) {
this.subscriptionFilters = subscriptionFilters;
}

@JsonAnyGetter
public Map<String, Object> getAdditionalProperties() {
return this.additionalProperties;
}

@JsonAnySetter
public void setAdditionalProperty(String name, Object value) {
this.additionalProperties.put(name, value);
}

@Override
public String toString() {
return ToStringBuilder.reflectionToString(this);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,11 @@

import java.io.IOException;

import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
import com.amazonaws.services.kinesis.model.Record;
import com.amazonaws.util.json.JSONArray;
import com.amazonaws.util.json.JSONException;
import com.amazonaws.util.json.JSONObject;
import com.sumologic.client.SimpleKinesisMessageModel;
import com.sumologic.client.implementations.SumologicTransformer;

Expand All @@ -13,18 +17,28 @@
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.nio.ByteBuffer;
import java.nio.CharBuffer;
import java.nio.charset.CharacterCodingException;
import java.nio.charset.Charset;
import java.nio.charset.CharsetDecoder;
import java.nio.charset.CharsetEncoder;
import java.util.List;
import java.util.zip.GZIPInputStream;

import com.fasterxml.jackson.databind.JsonMappingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.gson.Gson;

/**
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
* A custom transfomer for {@link CloudWatchLogsMessageModel} records in JSON format. The output is in a format
* usable for insertions to Sumologic.
*/
public class CloudWatchMessageModelSumologicTransformer implements
SumologicTransformer<SimpleKinesisMessageModel> {
public class CloudWatchMessageModelSumologicTransformer
implements SumologicTransformer<CloudWatchLogsMessageModel> {

private static final Log LOG = LogFactory.getLog(CloudWatchMessageModelSumologicTransformer.class);

private static CharsetEncoder encoder = Charset.forName("UTF-8").newEncoder();

/**
* Creates a new KinesisMessageModelSumologicTransformer.
Expand All @@ -34,24 +48,70 @@ public CloudWatchMessageModelSumologicTransformer() {
}

@Override
public String fromClass(SimpleKinesisMessageModel message) {
return message.toString();
}
public String fromClass(CloudWatchLogsMessageModel message) {
String jsonMessage = "";
JSONObject outputObject;

List<LogEvent> logEvents = message.getLogEvents();
int logEventsSize = logEvents.size();
for (int i=0;i<logEventsSize;i++) {
LogEvent log = logEvents.get(i);

outputObject = new JSONObject();
try {
// Header
outputObject.put("logGroup", message.getLogGroup());
outputObject.put("logStream", message.getLogStream());
outputObject.put("messageType", message.getMessageType());
outputObject.put("owner", message.getOwner());
outputObject.put("subscriptionFilters", new JSONArray(message.getSubscriptionFilters()));

// Body
outputObject.put("id", log.getId());
outputObject.put("message", log.getMessage());
outputObject.put("timestamp", log.getTimestamp());
} catch (JSONException e) {
LOG.error("Unable to convert message into JSON String: "+e.getMessage());
}
jsonMessage += outputObject.toString();
if (i < logEventsSize - 1) {
jsonMessage += '\n';
}
}
return jsonMessage;
}

@Override
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
public CloudWatchLogsMessageModel toClass(Record record) {
byte[] decodedRecord = record.getData().array();
String stringifiedRecord = decompressGzip(decodedRecord);

if (stringifiedRecord == null) {
LOG.error("Unable to decompress the record: "+new String(record.getData().array()));
LOG.error("Not attempting to transform into a Message Model");
LOG.error("Unable to decompress the record: "+new String(record.getData().array())
+"\nNot attempting to transform into a Message Model");
return null;
}

ByteBuffer bufferedData = null;
try {
bufferedData = encoder.encode(CharBuffer.wrap(stringifiedRecord));
} catch (CharacterCodingException e) {
LOG.error("Unable to set the decompressed Record for serializing "+e.getMessage());
}
record.setData(bufferedData);

return new SimpleKinesisMessageModel(stringifiedRecord);
try {
return new ObjectMapper().readValue(
record.getData().array(),
CloudWatchLogsMessageModel.class);
} catch (IOException e) {
LOG.error("Unable to convert the Record into a POJO: "+stringifiedRecord
+"\nerror: "+e.getMessage());
}
return null;

}

public static String decompressGzip(byte[] compressedData) {
try {
GZIPInputStream gis = new GZIPInputStream(new ByteArrayInputStream(compressedData));
Expand Down Expand Up @@ -83,4 +143,14 @@ public static String byteBufferToString(ByteBuffer buffer){
return data;
}

private static final Gson gson = new Gson();
public static boolean verifyJSON(String json) {
try {
gson.fromJson(json, Object.class);
return true;
} catch(com.google.gson.JsonSyntaxException ex) {
return false;
}
}

}
71 changes: 71 additions & 0 deletions src/main/java/com/sumologic/client/LogEvent.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
package com.sumologic.client;

import java.util.HashMap;
import java.util.Map;

import com.fasterxml.jackson.annotation.JsonAnyGetter;
import com.fasterxml.jackson.annotation.JsonAnySetter;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonPropertyOrder;

@JsonInclude(JsonInclude.Include.NON_NULL)
@JsonPropertyOrder({
"id",
"message",
"timestamp"
})

public class LogEvent {

@JsonProperty("id")
private String id;
@JsonProperty("message")
private String message;
@JsonProperty("timestamp")
private Long timestamp;
@JsonIgnore
private Map<String, Object> additionalProperties = new HashMap<String, Object>();

@JsonProperty("id")
public String getId() {
return id;
}

@JsonProperty("id")
public void setId(String id) {
this.id = id;
}

@JsonProperty("message")
public String getMessage() {
return message;
}

@JsonProperty("message")
public void setMessage(String message) {
this.message = message;
}

@JsonProperty("timestamp")
public Long getTimestamp() {
return timestamp;
}

@JsonProperty("timestamp")
public void setTimestamp(Long timestamp) {
this.timestamp = timestamp;
}

@JsonAnyGetter
public Map<String, Object> getAdditionalProperties() {
return this.additionalProperties;
}

@JsonAnySetter
public void setAdditionalProperty(String name, Object value) {
this.additionalProperties.put(name, value);
}

}
4 changes: 2 additions & 2 deletions src/main/java/com/sumologic/client/SumologicSender.java
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ public boolean sendToSumologic(String data) throws IOException{
BoundRequestBuilder builder = null;
builder = this.clientPreparePost(url);

byte[] compressedData = compressGzip(data);
byte[] compressedData = SumologicSender.compressGzip(data);

builder.setHeader("Content-Encoding", "gzip");
builder.setBody(compressedData);
Expand Down Expand Up @@ -82,7 +82,7 @@ public boolean sendToSumologic(String data) throws IOException{
}
}

public byte[] compressGzip(String data) {
public static byte[] compressGzip(String data) {
if (data == null || data.length() == 0) {
return null;
}
Expand Down
Loading