Skip to content

Commit 4c8b07d

Browse files
Juan Pablo Diaz-VazJuan Pablo Diaz-Vaz
authored andcommitted
SUMOK-18 Changed name of MessageModel class
1 parent 6bc61b3 commit 4c8b07d

9 files changed

Lines changed: 61 additions & 41 deletions

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

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@
2828
import org.apache.commons.logging.Log;
2929
import org.apache.commons.logging.LogFactory;
3030

31-
import com.mcplusa.sumologic.KinesisMessageModel;
31+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
3232

3333
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
3434
import com.amazonaws.services.kinesis.model.PutRecordRequest;
@@ -41,15 +41,15 @@ public class BatchedStreamSource extends StreamSource {
4141
private static Log LOG = LogFactory.getLog(BatchedStreamSource.class);
4242

4343
private static int NUM_BYTES_PER_PUT_REQUEST = 50000;
44-
List<KinesisMessageModel> buffer;
44+
List<SimpleKinesisMessageModel> buffer;
4545

4646
public BatchedStreamSource(KinesisConnectorConfiguration config, String inputFile) {
4747
this(config, inputFile, false);
4848
}
4949

5050
public BatchedStreamSource(KinesisConnectorConfiguration config, String inputFile, boolean loopOverStreamSource) {
5151
super(config, inputFile, loopOverStreamSource);
52-
buffer = new ArrayList<KinesisMessageModel>();
52+
buffer = new ArrayList<SimpleKinesisMessageModel>();
5353
}
5454

5555
@Override
@@ -59,14 +59,14 @@ protected void processInputStream(InputStream inputStream, int iteration) throws
5959
int lines = 0;
6060

6161
while ((line = br.readLine()) != null) {
62-
KinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, KinesisMessageModel.class);
62+
SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
6363
buffer.add(kinesisMessageModel);
6464
if (numBytesInBuffer() > NUM_BYTES_PER_PUT_REQUEST) {
6565
/*
6666
* We need to remove the last record to ensure this data blob is accepted by the Amazon Kinesis
6767
* client which restricts the data blob to be less than 50 KB.
6868
*/
69-
KinesisMessageModel lastRecord = buffer.remove(buffer.size() - 1);
69+
SimpleKinesisMessageModel lastRecord = buffer.remove(buffer.size() - 1);
7070
flushBuffer();
7171
/*
7272
* We add it back so it will be part of the next batch.

src/main/java/com/mcplusa/kinesis/StreamSource.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@
2323
import org.apache.commons.logging.Log;
2424
import org.apache.commons.logging.LogFactory;
2525

26-
import com.mcplusa.sumologic.KinesisMessageModel;
26+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
2727
import com.mcplusa.kinesis.utils.KinesisUtils;
2828

2929
import com.amazonaws.auth.AWSCredentialsProvider;
@@ -120,7 +120,7 @@ protected void processInputStream(InputStream inputStream, int iteration) throws
120120
String line;
121121
int lines = 0;
122122
while ((line = br.readLine()) != null) {
123-
KinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, KinesisMessageModel.class);
123+
SimpleKinesisMessageModel kinesisMessageModel = objectMapper.readValue(line, SimpleKinesisMessageModel.class);
124124

125125
PutRecordRequest putRecordRequest = new PutRecordRequest();
126126
putRecordRequest.setStreamName(config.KINESIS_INPUT_STREAM);

src/main/java/com/mcplusa/sumologic/KinesisMessageModelSumologicTransformer.java renamed to src/main/java/com/mcplusa/sumologic/CloudWatchMessageModelSumologicTransformer.java

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44

55
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
66
import com.amazonaws.services.kinesis.model.Record;
7-
import com.mcplusa.sumologic.KinesisMessageModel;
7+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
88
import com.mcplusa.sumologic.implementations.SumologicEmitter;
99
import com.mcplusa.sumologic.implementations.SumologicTransformer;
1010

@@ -22,32 +22,32 @@
2222

2323

2424
/**
25-
* A custom transfomer for {@link KinesisMessageModel} records in JSON format. The output is in a format
25+
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
2626
* usable for insertions to Sumologic.
2727
*/
28-
public class KinesisMessageModelSumologicTransformer implements
29-
SumologicTransformer<KinesisMessageModel> {
28+
public class CloudWatchMessageModelSumologicTransformer implements
29+
SumologicTransformer<SimpleKinesisMessageModel> {
3030

31-
private static final Log LOG = LogFactory.getLog(KinesisMessageModelSumologicTransformer.class);
31+
private static final Log LOG = LogFactory.getLog(CloudWatchMessageModelSumologicTransformer.class);
3232

3333
/**
3434
* Creates a new KinesisMessageModelSumologicTransformer.
3535
*/
36-
public KinesisMessageModelSumologicTransformer() {
36+
public CloudWatchMessageModelSumologicTransformer() {
3737
super();
3838
}
3939

4040
@Override
41-
public String fromClass(KinesisMessageModel message) {
41+
public String fromClass(SimpleKinesisMessageModel message) {
4242
return message.toString();
4343
}
4444

4545
@Override
46-
public KinesisMessageModel toClass(Record record) throws IOException {
46+
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
4747
byte[] decodedRecord = record.getData().array();
4848
String stringifiedRecord = decompressGzip(decodedRecord);
4949

50-
return new KinesisMessageModel(stringifiedRecord);
50+
return new SimpleKinesisMessageModel(stringifiedRecord);
5151
}
5252

5353
public static String decompressGzip(byte[] compressedData) {

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

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@
44

55
import com.amazonaws.services.kinesis.connectors.BasicJsonTransformer;
66
import com.amazonaws.services.kinesis.model.Record;
7-
import com.mcplusa.sumologic.KinesisMessageModel;
7+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
88
import com.mcplusa.sumologic.implementations.SumologicEmitter;
99
import com.mcplusa.sumologic.implementations.SumologicTransformer;
1010

@@ -22,11 +22,11 @@
2222

2323

2424
/**
25-
* A custom transfomer for {@link KinesisMessageModel} records in JSON format. The output is in a format
25+
* A custom transfomer for {@link SimpleKinesisMessageModel} records in JSON format. The output is in a format
2626
* usable for insertions to Sumologic.
2727
*/
2828
public class DefaultKinesisMessageModelSumologicTransformer implements
29-
SumologicTransformer<KinesisMessageModel> {
29+
SumologicTransformer<SimpleKinesisMessageModel> {
3030

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

@@ -38,15 +38,15 @@ public DefaultKinesisMessageModelSumologicTransformer() {
3838
}
3939

4040
@Override
41-
public String fromClass(KinesisMessageModel message) {
41+
public String fromClass(SimpleKinesisMessageModel message) {
4242
return message.toString();
4343
}
4444

4545
@Override
46-
public KinesisMessageModel toClass(Record record) throws IOException {
46+
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
4747
byte[] decodedRecord = record.getData().array();
4848
String stringifiedRecord = new String(decodedRecord);
4949

50-
return new KinesisMessageModel(stringifiedRecord);
50+
return new SimpleKinesisMessageModel(stringifiedRecord);
5151
}
5252
}

src/main/java/com/mcplusa/sumologic/KinesisMessageModel.java renamed to src/main/java/com/mcplusa/sumologic/SimpleKinesisMessageModel.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,11 @@
22

33
import java.io.Serializable;
44

5-
public class KinesisMessageModel implements Serializable {
5+
public class SimpleKinesisMessageModel implements Serializable {
66
private String data;
77
private int id;
88

9-
public KinesisMessageModel(String data) {
9+
public SimpleKinesisMessageModel(String data) {
1010
this.data = data;
1111
this.id = 1;
1212
}

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

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,10 @@
33
import com.amazonaws.services.kinesis.connectors.KinesisConnectorRecordProcessorFactory;
44

55
import com.mcplusa.kinesis.KinesisConnectorExecutor;
6-
import com.mcplusa.sumologic.KinesisMessageModel;
6+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
77
import com.mcplusa.sumologic.SumologicMessageModelPipeline;
88

9-
public class SumologicExecutor extends KinesisConnectorExecutor<KinesisMessageModel, String> {
9+
public class SumologicExecutor extends KinesisConnectorExecutor<SimpleKinesisMessageModel, String> {
1010

1111
private static String configFile = "SumologicConnector.properties";
1212

@@ -19,9 +19,9 @@ public SumologicExecutor(String configFile) {
1919
}
2020

2121
@Override
22-
public KinesisConnectorRecordProcessorFactory<KinesisMessageModel, String>
22+
public KinesisConnectorRecordProcessorFactory<SimpleKinesisMessageModel, String>
2323
getKinesisConnectorRecordProcessorFactory() {
24-
return new KinesisConnectorRecordProcessorFactory<KinesisMessageModel, String>
24+
return new KinesisConnectorRecordProcessorFactory<SimpleKinesisMessageModel, String>
2525
(new SumologicMessageModelPipeline(),config);
2626
}
2727

@@ -30,7 +30,7 @@ public SumologicExecutor(String configFile) {
3030
* @param args
3131
*/
3232
public static void main(String[] args) {
33-
KinesisConnectorExecutor<KinesisMessageModel, String> sumologicExecutor =
33+
KinesisConnectorExecutor<SimpleKinesisMessageModel, String> sumologicExecutor =
3434
new SumologicExecutor(configFile);
3535
sumologicExecutor.run();
3636
}

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

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

6-
import com.mcplusa.sumologic.KinesisMessageModel;
7-
import com.mcplusa.sumologic.KinesisMessageModelSumologicTransformer;
6+
import com.mcplusa.sumologic.SimpleKinesisMessageModel;
7+
import com.mcplusa.sumologic.CloudWatchMessageModelSumologicTransformer;
88
import com.mcplusa.sumologic.implementations.SumologicEmitter;
99
import com.amazonaws.services.kinesis.connectors.interfaces.IKinesisConnectorPipeline;
1010
import com.amazonaws.services.kinesis.connectors.KinesisConnectorConfiguration;
@@ -22,12 +22,12 @@
2222
* <ul>
2323
* <li>{@link SumologicEmitter}</li>
2424
* <li>{@link BasicMemoryBuffer}</li>
25-
* <li>{@link KinesisMessageModelSumologicTransformer}</li>
25+
* <li>{@link CloudWatchMessageModelSumologicTransformer}</li>
2626
* <li>{@link AllPassFilter}</li>
2727
* </ul>
2828
*/
2929
public class SumologicMessageModelPipeline implements
30-
IKinesisConnectorPipeline<KinesisMessageModel, String> {
30+
IKinesisConnectorPipeline<SimpleKinesisMessageModel, String> {
3131

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

@@ -37,12 +37,12 @@ public IEmitter<String> getEmitter(KinesisConnectorConfiguration configuration)
3737
}
3838

3939
@Override
40-
public IBuffer<KinesisMessageModel> getBuffer(KinesisConnectorConfiguration configuration) {
41-
return new BasicMemoryBuffer<KinesisMessageModel>(configuration);
40+
public IBuffer<SimpleKinesisMessageModel> getBuffer(KinesisConnectorConfiguration configuration) {
41+
return new BasicMemoryBuffer<SimpleKinesisMessageModel>(configuration);
4242
}
4343

4444
@Override
45-
public ITransformer<KinesisMessageModel, String>
45+
public ITransformer<SimpleKinesisMessageModel, String>
4646
getTransformer(KinesisConnectorConfiguration configuration) {
4747

4848
// Load specified class
@@ -52,7 +52,7 @@ public IBuffer<KinesisMessageModel> getBuffer(KinesisConnectorConfiguration conf
5252
Class ModelClass = null;
5353
try {
5454
ModelClass = classLoader.loadClass(className);
55-
ITransformer<KinesisMessageModel, String> ITransformerObject = (ITransformer<KinesisMessageModel, String>)ModelClass.newInstance();
55+
ITransformer<SimpleKinesisMessageModel, String> ITransformerObject = (ITransformer<SimpleKinesisMessageModel, String>)ModelClass.newInstance();
5656
LOG.info("Using transformer: "+ITransformerObject.getClass().getName());
5757
return ITransformerObject;
5858
} catch (ClassNotFoundException e) {
@@ -67,8 +67,8 @@ public IBuffer<KinesisMessageModel> getBuffer(KinesisConnectorConfiguration conf
6767
}
6868

6969
@Override
70-
public IFilter<KinesisMessageModel> getFilter(KinesisConnectorConfiguration configuration) {
71-
return new AllPassFilter<KinesisMessageModel>();
70+
public IFilter<SimpleKinesisMessageModel> getFilter(KinesisConnectorConfiguration configuration) {
71+
return new AllPassFilter<SimpleKinesisMessageModel>();
7272
}
7373

7474
}

src/main/java/com/mcplusa/sumologic/SumologicSender.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ public boolean sendToSumologic(String data) throws IOException{
4343
int statusCode;
4444

4545
do {
46-
HttpPost post = null;
46+
HttpPost post = null;
4747
post = new HttpPost(url);
4848
post.setEntity(new StringEntity(data, HTTP.PLAIN_TEXT_TYPE, HTTP.UTF_8));
4949
HttpResponse response = httpClient.execute(post);

src/main/java/com/mcplusa/sumologic/implementations/SumologicEmitter.java

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ public class SumologicEmitter implements IEmitter<String> {
2626

2727
private SumologicSender sender;
2828
private KinesisConnectorForSumologicConfiguration config;
29+
private static final boolean SEND_RECORDS_IN_BATCHES = true;
2930

3031
public SumologicEmitter(KinesisConnectorConfiguration configuration) {
3132
this.config = (KinesisConnectorForSumologicConfiguration) configuration;
@@ -40,7 +41,11 @@ public SumologicEmitter(String url) {
4041
public List<String> emit(final UnmodifiableBuffer<String> buffer)
4142
throws IOException {
4243
List<String> records = buffer.getRecords();
43-
return sendBatchConcatenating(records);
44+
if (SEND_RECORDS_IN_BATCHES) {
45+
return sendBatchConcatenating(records);
46+
} else {
47+
return sendRecordsOneByOne(records);
48+
}
4449
}
4550

4651
public List<String> sendBatchConcatenating(List<String> records) {
@@ -65,6 +70,21 @@ public List<String> sendBatchConcatenating(List<String> records) {
6570
return records;
6671
}
6772
}
73+
74+
public List<String> sendRecordsOneByOne (List<String> records) {
75+
ArrayList<String> failedRecords = new ArrayList<String>();
76+
for (String record: records) {
77+
try {
78+
if (!sender.sendToSumologic(record)) {
79+
failedRecords.add(record);
80+
}
81+
} catch (IOException e) {
82+
LOG.warn("Couldn't send record: "+record);
83+
}
84+
}
85+
LOG.info("Sent records: "+(records.size()-failedRecords.size())+" failed: "+failedRecords.size());
86+
return failedRecords;
87+
}
6888

6989
@Override
7090
public void fail(List<String> records) {

0 commit comments

Comments
 (0)