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 @@ -34,6 +34,7 @@
<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= "http://central.maven.org/maven2/log4j/log4j/1.2.17/log4j-1.2.17.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
10 changes: 0 additions & 10 deletions log4j.properties
Original file line number Diff line number Diff line change
@@ -1,16 +1,6 @@
# Root logger option
log4j.rootLogger=INFO, stdout
log4j.logger.sumologic = TRACE, sumo

# Direct log messages to sumo
log4j.appender.sumo=com.sumologic.log4j.BufferedSumoLogicAppender
log4j.appender.sumo.url=https://collectors.us2.sumologic.com/receiver/v1/http/ZaVnC4dhaV0GzIY4tZaKLL26afV52gXBvSFc3jG1eLc2lKINzS2doZdRjUMQMb2CXK8r6fdmHoUazJiHjJ-2OygApoWFaxCTkWFrzAiraCc5i411pkio-g==
log4j.appender.sumo.layout=org.apache.log4j.PatternLayout
log4j.appender.sumo.layout.ConversionPattern=%d{DATE} %5p %c{1}:%L - %m%n
log4j.additivity.sumo = false
log4j.appender.sumo.Threshold = TRACE

log4j.appender.stdout=org.apache.log4j.ConsoleAppender
log4j.appender.stdout.layout=org.apache.log4j.PatternLayout
log4j.appender.stdout.layout.ConversionPattern=%d{DATE} %5p %c{1}:%L - %m%n
log4j.appender.stdout.Threshold = INFO
Original file line number Diff line number Diff line change
Expand Up @@ -2,42 +2,33 @@

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;
import com.sumologic.client.model.CloudWatchLogsMessageModel;
import com.sumologic.client.model.LogEvent;

import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.log4j.Logger;

import java.io.ByteArrayInputStream;
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 CloudWatchLogsMessageModel} records in JSON format. The output is in a format
* usable for insertions to Sumologic.
*/
public class CloudWatchMessageModelSumologicTransformer
implements SumologicTransformer<CloudWatchLogsMessageModel> {

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

private static final Logger LOG = Logger.getLogger(CloudWatchMessageModelSumologicTransformer.class.getName());

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

/**
Expand Down Expand Up @@ -84,7 +75,7 @@ public String fromClass(CloudWatchLogsMessageModel message) {
@Override
public CloudWatchLogsMessageModel toClass(Record record) {
byte[] decodedRecord = record.getData().array();
String stringifiedRecord = decompressGzip(decodedRecord);
String stringifiedRecord = SumologicKinesisUtils.decompressGzip(decodedRecord);

if (stringifiedRecord == null) {
LOG.error("Unable to decompress the record: "+new String(record.getData().array())
Expand All @@ -109,48 +100,5 @@ public CloudWatchLogsMessageModel toClass(Record record) {
+"\nerror: "+e.getMessage());
}
return null;

}

public static String decompressGzip(byte[] compressedData) {
try {
GZIPInputStream gis = new GZIPInputStream(new ByteArrayInputStream(compressedData));
BufferedReader bf = new BufferedReader(new InputStreamReader(gis, "UTF-8"));

String outStr = "";
String line;
while ((line=bf.readLine())!=null) {
outStr += line;
}
return outStr;
} catch (IOException exc) {
LOG.warn("Exception during decompression of data: " + exc.getMessage());
return null;
}
}

public static String byteBufferToString(ByteBuffer buffer){
String data = "";
CharsetDecoder decoder = Charset.forName("UTF-8").newDecoder();
try{
int old_position = buffer.position();
data = decoder.decode(buffer).toString();
buffer.position(old_position);
}catch (Exception e){
e.printStackTrace();
return "";
}
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;
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -2,34 +2,16 @@

import java.io.IOException;

import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
import com.amazonaws.services.kinesis.model.Record;
import com.sumologic.client.SimpleKinesisMessageModel;
import com.sumologic.client.implementations.SumologicEmitter;
import com.sumologic.client.implementations.SumologicTransformer;

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

import java.io.ByteArrayInputStream;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import java.util.zip.GZIPInputStream;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;

import org.apache.commons.codec.binary.Base64;

import com.sumologic.client.model.SimpleKinesisMessageModel;

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

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

/**
* Creates a new KinesisMessageModelSumologicTransformer.
*/
Expand All @@ -46,7 +28,7 @@ public String fromClass(SimpleKinesisMessageModel message) {
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
byte[] decodedRecord = record.getData().array();
String stringifiedRecord = new String(decodedRecord);

return new SimpleKinesisMessageModel(stringifiedRecord);
}
}
6 changes: 2 additions & 4 deletions src/main/java/com/sumologic/client/SumologicExecutor.java
Original file line number Diff line number Diff line change
@@ -1,13 +1,11 @@
package com.sumologic.client;

import com.amazonaws.services.kinesis.connectors.KinesisConnectorRecordProcessorFactory;

import com.sumologic.kinesis.KinesisConnectorRecordProcessorFactory;
import com.sumologic.kinesis.KinesisConnectorExecutor;
import com.sumologic.client.SimpleKinesisMessageModel;
import com.sumologic.client.SumologicMessageModelPipeline;
import com.sumologic.client.model.SimpleKinesisMessageModel;

public class SumologicExecutor extends KinesisConnectorExecutor<SimpleKinesisMessageModel, String> {

private static String configFile = "SumologicConnector.properties";

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

import java.io.BufferedReader;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStreamReader;
import java.nio.ByteBuffer;
import java.nio.charset.Charset;
import java.nio.charset.CharsetDecoder;
import java.util.zip.GZIPInputStream;
import java.util.zip.GZIPOutputStream;

import org.apache.log4j.Logger;

import com.google.gson.Gson;

public class SumologicKinesisUtils {
private static final Logger LOG = Logger.getLogger(SumologicKinesisUtils.class.getName());

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

ByteArrayOutputStream outputStream=new ByteArrayOutputStream();
GZIPOutputStream gzip;
try {
gzip = new GZIPOutputStream(outputStream);
} catch (IOException e) {
LOG.error("Cannot compress into GZIP "+e.getMessage());
return null;
}

// Put data into the GZIP buffer
try {
gzip.write(data.getBytes("UTF-8"));
gzip.close();
} catch (IOException e) {
e.printStackTrace();
}

return outputStream.toByteArray();
}

public static String decompressGzip(byte[] compressedData) {
try {
GZIPInputStream gis = new GZIPInputStream(new ByteArrayInputStream(compressedData));
BufferedReader bf = new BufferedReader(new InputStreamReader(gis, "UTF-8"));

String outStr = "";
String line;
while ((line=bf.readLine())!=null) {
outStr += line;
}
return outStr;
} catch (IOException exc) {
LOG.warn("Exception during decompression of data: " + exc.getMessage());
return null;
}
}

public static String byteBufferToString(ByteBuffer buffer){
String data = "";
CharsetDecoder decoder = Charset.forName("UTF-8").newDecoder();
try{
int old_position = buffer.position();
data = decoder.decode(buffer).toString();
buffer.position(old_position);
}catch (Exception e){
e.printStackTrace();
return "";
}
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;
}
}
}
Original file line number Diff line number Diff line change
@@ -1,11 +1,9 @@
package com.sumologic.client;

import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.log4j.Logger;

import com.sumologic.client.SimpleKinesisMessageModel;
import com.sumologic.client.CloudWatchMessageModelSumologicTransformer;
import com.sumologic.client.implementations.SumologicEmitter;
import com.sumologic.client.model.SimpleKinesisMessageModel;
import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
import com.amazonaws.services.kinesis.connectors.impl.BasicMemoryBuffer;
Expand All @@ -28,9 +26,8 @@
*/
public class SumologicMessageModelPipeline implements
IKinesisConnectorPipeline<SimpleKinesisMessageModel, String> {
private static final Logger LOG = Logger.getLogger(SumologicMessageModelPipeline.class.getName());

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

@Override
public IEmitter<String> getEmitter(KinesisConnectorConfiguration configuration) {
return new SumologicEmitter(configuration);
Expand Down
27 changes: 1 addition & 26 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 = SumologicSender.compressGzip(data);
byte[] compressedData = SumologicKinesisUtils.compressGzip(data);

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

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

ByteArrayOutputStream outputStream=new ByteArrayOutputStream();
GZIPOutputStream gzip;
try {
gzip = new GZIPOutputStream(outputStream);
} catch (IOException e) {
LOG.error("Cannot compress into GZIP "+e.getMessage());
return null;
}

// Put data into the GZIP buffer
try {
gzip.write(data.getBytes("UTF-8"));
gzip.close();
} catch (IOException e) {
e.printStackTrace();
}

return outputStream.toByteArray();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,12 @@

import java.io.IOException;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;

import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.log4j.Logger;

import com.sumologic.client.SumologicSender;
import com.sumologic.client.KinesisConnectorForSumologicConfiguration;

import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
import com.amazonaws.services.kinesis.connectors.UnmodifiableBuffer;
import com.amazonaws.services.kinesis.connectors.interfaces.IEmitter;
Expand All @@ -22,7 +18,7 @@
* Sumologic.
*/
public class SumologicEmitter implements IEmitter<String> {
private static final Log LOG = LogFactory.getLog(SumologicEmitter.class);
private static final Logger LOG = Logger.getLogger(SumologicEmitter.class.getName());

private SumologicSender sender;
private KinesisConnectorForSumologicConfiguration config;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package com.sumologic.client;
package com.sumologic.client.model;

import java.util.ArrayList;
import java.util.HashMap;
Expand All @@ -7,7 +7,6 @@

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;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package com.sumologic.client;
package com.sumologic.client.model;

import java.util.HashMap;
import java.util.Map;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package com.sumologic.client;
package com.sumologic.client.model;

import java.io.Serializable;

Expand Down
Loading