-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathDefaultKinesisMessageModelSumologicTransformer.java
More file actions
52 lines (40 loc) · 1.63 KB
/
Copy pathDefaultKinesisMessageModelSumologicTransformer.java
File metadata and controls
52 lines (40 loc) · 1.63 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
package com.sumologic.client;
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;
/**
* 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.
*/
public DefaultKinesisMessageModelSumologicTransformer() {
super();
}
@Override
public String fromClass(SimpleKinesisMessageModel message) {
return message.toString();
}
@Override
public SimpleKinesisMessageModel toClass(Record record) throws IOException {
byte[] decodedRecord = record.getData().array();
String stringifiedRecord = new String(decodedRecord);
return new SimpleKinesisMessageModel(stringifiedRecord);
}
}