-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathSumologicMessageModelPipeline.java
More file actions
71 lines (60 loc) · 3.04 KB
/
Copy pathSumologicMessageModelPipeline.java
File metadata and controls
71 lines (60 loc) · 3.04 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
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
package com.sumologic.client;
import org.apache.log4j.Logger;
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;
import com.amazonaws.services.kinesis.connectors.impl.AllPassFilter;
import com.amazonaws.services.kinesis.connectors.interfaces.IEmitter;
import com.amazonaws.services.kinesis.connectors.interfaces.IBuffer;
import com.amazonaws.services.kinesis.connectors.interfaces.ITransformer;
import com.amazonaws.services.kinesis.connectors.interfaces.IFilter;
/**
* The Pipeline used by the Sumologic. Processes KinesisMessageModel records in JSON String
* format. Uses:
* <ul>
* <li>{@link SumologicEmitter}</li>
* <li>{@link BasicMemoryBuffer}</li>
* <li>{@link CloudWatchMessageModelSumologicTransformer}</li>
* <li>{@link AllPassFilter}</li>
* </ul>
*/
public class SumologicMessageModelPipeline implements
IKinesisConnectorPipeline<SimpleKinesisMessageModel, String> {
private static final Logger LOG = Logger.getLogger(SumologicMessageModelPipeline.class.getName());
@Override
public IEmitter<String> getEmitter(KinesisConnectorConfiguration configuration) {
return new SumologicEmitter(configuration);
}
@Override
public IBuffer<SimpleKinesisMessageModel> getBuffer(KinesisConnectorConfiguration configuration) {
return new BasicMemoryBuffer<SimpleKinesisMessageModel>(configuration);
}
@Override
public ITransformer<SimpleKinesisMessageModel, String>
getTransformer(KinesisConnectorConfiguration configuration) {
// Load specified class
String argClass = ((KinesisConnectorForSumologicConfiguration)configuration).TRANSFORMER_CLASS;
String className = "com.sumologic.client."+argClass;
ClassLoader classLoader = SumologicMessageModelPipeline.class.getClassLoader();
Class ModelClass = null;
try {
ModelClass = classLoader.loadClass(className);
ITransformer<SimpleKinesisMessageModel, String> ITransformerObject = (ITransformer<SimpleKinesisMessageModel, String>)ModelClass.newInstance();
LOG.info("Using transformer: "+ITransformerObject.getClass().getName());
return ITransformerObject;
} catch (ClassNotFoundException e) {
LOG.error("Class not found: "+className+" error: "+e.getMessage());
} catch (InstantiationException e) {
LOG.error("Class not found: "+className+" error: "+e.getMessage());
} catch (IllegalAccessException e) {
LOG.error("Class not found: "+className+" error: "+e.getMessage());
}
return new DefaultKinesisMessageModelSumologicTransformer();
}
@Override
public IFilter<SimpleKinesisMessageModel> getFilter(KinesisConnectorConfiguration configuration) {
return new AllPassFilter<SimpleKinesisMessageModel>();
}
}